Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
0fe9bb0
feat: add MongoDB ObjectID filtering support and corresponding test
flarco Sep 4, 2025
5561e9f
feat: refactor column retrieval and validation in CopyViaStage method
flarco Sep 5, 2025
9308044
feat: initialize sample size from environment variable in SetSampleSi…
flarco Sep 6, 2025
af590f2
feat: add type checking for new columns in parseRecords method
flarco Sep 6, 2025
f6a6836
feat: add JSON to Snowflake nested import test and corresponding YAML…
flarco Sep 6, 2025
25329ea
feat: enhance DropTable method to attempt drop after purge failure
flarco Sep 6, 2025
e476203
feat: improve error message for unexpected mapping structure in GetTa…
flarco Sep 6, 2025
a99b76f
feat: update GetTableColumns and StreamRowsContext methods to use 'an…
flarco Sep 6, 2025
ca5e40b
refactor: move evaluator initialization outside of test loop in TestE…
flarco Sep 6, 2025
84f4da8
feat: add version handling to ClickhouseConn and update connection logic
flarco Sep 6, 2025
bbb42d3
fix: correct SQL delete statement syntax in GenerateMergeSQL method
flarco Sep 6, 2025
4f9e864
fix: update internal stage handling to require fully qualified stage …
flarco Sep 6, 2025
ede70dc
fix: clarify warning message for fully qualified stage name requirement
flarco Sep 6, 2025
729968d
feat: add Merge method to ClickhouseConn for handling upserts with re…
flarco Sep 6, 2025
955f244
fix: update error messages to use 'merge' terminology in Merge functions
flarco Sep 6, 2025
59484e2
fix: update LoadAPISpec to use specIdentifier and add FetchSpec for f…
flarco Sep 8, 2025
9669a03
fix: enhance GenerateMergeConfig to handle primary key scenarios in s…
flarco Sep 8, 2025
ff5e541
test: add validation check for email hash length in staged transforms
flarco Sep 8, 2025
544c74d
fix: update Prometheus YAML configurations to ensure proper date form…
flarco Sep 8, 2025
840e59c
fix: update error messages in ErrorHelper for clarity and improved do…
flarco Sep 9, 2025
18cfd57
Revert prometheus changes
flarco Sep 9, 2025
97875b0
fix: update transform field names to use PascalCase for consistency
flarco Sep 9, 2025
3d85691
feat: enhance Prometheus data handling with transform support and con…
flarco Sep 9, 2025
bd6ef47
fix: add fields parameter to SQL queries in replication and task exec…
flarco Sep 11, 2025
f48dd98
fix: add statement_timeout support in connection URL and enhance erro…
flarco Sep 11, 2025
26f442b
feat: add create_schema template to Databricks YAML configuration
flarco Sep 12, 2025
10b1994
fix: update internal volume handling to store and retrieve volume pat…
flarco Sep 12, 2025
295232a
feat: add test for SQL Server money type to StarRocks replication
flarco Sep 12, 2025
5809100
fix: add strict parameter to StreamLoad for enhanced data validation
flarco Sep 12, 2025
c82e1b7
fix: enforce minimum precision of 15 for SQL column types to address …
flarco Sep 13, 2025
4d74d0c
feat: enhance SQL replication tests by increasing decimal precision a…
flarco Sep 13, 2025
3e74e24
fix: update jinzhu/copier dependency to remove indirect requirement
flarco Sep 13, 2025
92e0fe0
feat: enhance Redshift credential handling by consolidating AWS crede…
flarco Sep 13, 2025
88db8a5
fix: update minimum precision check for SQL columns to use configurab…
flarco Sep 14, 2025
09aae46
fix: remove redundant streams entry in MSSQL test configuration
flarco Sep 14, 2025
42e5fdc
feat: add YellowString function for colored output in terminal
flarco Sep 14, 2025
a7d7470
feat: update Prometheus buffer configuration to use parquet format an…
flarco Sep 14, 2025
3566a1b
feat: add Casted method to Transform interface and enhance type handl…
flarco Sep 14, 2025
ba10bb3
fix: ensure col.Sourced is set to false when conditions are not met i…
flarco Sep 14, 2025
49a7183
fix: remove redundant Stringer case in CastVal and update CastToStrin…
flarco Sep 14, 2025
8b7c91a
fix: update precision and scale handling for DecimalType across datab…
flarco Sep 15, 2025
0ef315f
fix: enhance decimal precision and scale handling in ColumnTyping and…
flarco Sep 15, 2025
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions cmd/sling/tests/replications/r.25.prometheus_buffer.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,14 @@ target: local
defaults:
mode: full-refresh
target_options:
format: csv
format: parquet
file_max_rows: 10000

streams:
prometheus_metrics:
transforms:
- '*': cast(value, "string")
- event_date: date_format(now(), "%Y-%m-%d")
sql: |
up#{"start": "now-30m", "end": "now", "step": "5m"}
object: '/tmp/prometheus_test_{run_timestamp}.csv'
sum(go_gc_duration_seconds) by (job, instance, quantile) # {"start": "now-2M"}
object: '/tmp/output/prometheus_test.parquet'
8 changes: 8 additions & 0 deletions cmd/sling/tests/replications/r.44.staged_transforms_test.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ streams:
- "*": cast(value, "string")
- "*": remove_diacritics(value)
- calculated_amount: 'cast(value, "decimal")'
event_date: date_format(now(), "%Y-%m-%d")

hooks:
post:
Expand All @@ -102,6 +103,7 @@ streams:
full_name,
email_hash,
calculated_amount,
event_date,
category,
display_name,
processed_status
Expand All @@ -126,6 +128,12 @@ streams:
- type: check
check: store.result[0].email == "john@example.com" # lowercased

- type: check
check: length(store.result[0].email_hash) > 30

- type: check
check: length(store.result[0].event_date) == 10

- type: check
check: store.result[0].full_name == "John Doe" # record reference combining trimmed names

Expand Down
39 changes: 39 additions & 0 deletions cmd/sling/tests/replications/r.54.mongo_objectid_filter.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
source: mongo
target: postgres

defaults:
mode: full-refresh
source_options:
flatten: 1

hooks:
start:
# Clean up test table if exists
- type: query
connection: '{target.name}'
query: DROP TABLE IF EXISTS public.mongo_objectid_test;

end:
# Verify records were inserted
- type: query
connection: '{target.name}'
query: SELECT COUNT(*) as count FROM public.mongo_objectid_test
into: result

- type: check
check: int_parse(store.result[0].count) == 90
failure_message: Should have 90 rows

# Clean up test table
- type: query
connection: '{target.name}'
query: DROP TABLE IF EXISTS public.mongo_objectid_test;

streams:
# Test ObjectID filtering on existing collection
# This tests that ObjectID hex strings in filters are properly converted
default.test1k_mongodb:
object: public.mongo_objectid_test
# Filter using ObjectID range - these will be converted to proper ObjectID types
# Using a range that should capture some documents
where: '{"_id": {"$gte": "67859d8ee682ab32317abc6f", "$lte": "67859d8ee682ab32317abcc8"}}'
41 changes: 41 additions & 0 deletions cmd/sling/tests/replications/r.55.json_snowflake_nested.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
source: LOCAL
target: SNOWFLAKE

hooks:
start:
- type: command
command: bash cmd/sling/tests/replications/r.55.json_snowflake_nested.yaml.prep.sh

- type: query
connection: '{target.name}'
query: DROP TABLE IF EXISTS public.sling_nested_json;

end:
# Verify total records were imported
- type: query
connection: '{target.name}'
query: SELECT COUNT(*) as count FROM public.sling_nested_json
into: result

- type: check
check: int_parse(store.result[0].count) == 1700
success_message: "✓ All 1700 JSON records imported successfully (1200 from nested1.json + 500 from nested2.json)"

# Cleanup files
- type: command
command: rm -rf /tmp/sling-test-files/

# Cleanup table
- type: query
connection: '{target.name}'
query: DROP TABLE IF EXISTS public.sling_nested_json;

streams:
file:///tmp/sling-test-files/*.json:
object: public.sling_nested_json
mode: full-refresh
source_options:
flatten: 1

env:
SLING_SAMPLE_SIZE: 100
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
mkdir -p /tmp/sling-test-files

# Create nested1.json with 1200 records (2 columns each)
echo '[' > /tmp/sling-test-files/nested1.json
for i in {1..1200}; do
if [ $i -gt 1 ]; then
echo "," >> /tmp/sling-test-files/nested1.json
fi
if [ $i -gt 1100 ]; then
echo -n "{\"id\": $i, \"name\": \"User$i\", \"email\": \"user$(($i))@example.com\"}" >> /tmp/sling-test-files/nested1.json
else
echo -n "{\"id\": $i, \"name\": \"User$i\"}" >> /tmp/sling-test-files/nested1.json
fi
done
echo -e "\n]" >> /tmp/sling-test-files/nested1.json

# Create nested2.json with 500 records (3 columns each)
echo '[' > /tmp/sling-test-files/nested2.json
for i in {1..500}; do
if [ $i -gt 1 ]; then
echo "," >> /tmp/sling-test-files/nested2.json
fi
value=$((i * 100))
echo -n "{\"id\": $((1200 + i)), \"email\": \"user$((1200 + i))@example.com\", \"value\": $value}" >> /tmp/sling-test-files/nested2.json
done
echo -e "\n]" >> /tmp/sling-test-files/nested2.json
220 changes: 220 additions & 0 deletions cmd/sling/tests/replications/r.56.mssql_starrocks_money.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
source: mssql
target: starrocks

defaults:
mode: full-refresh

hooks:
start:
- type: query
connection: '{source.name}'
query: |
-- Drop temp table if it exists
IF OBJECT_ID('dbo.money_test', 'U') IS NOT NULL DROP TABLE dbo.money_test;

-- Create table with money and decimal columns
CREATE TABLE dbo.money_test (
id INT PRIMARY KEY,
description VARCHAR(100),
money_col MONEY,
smallmoney_col SMALLMONEY,
decimal_col DECIMAL(24,6),
decimal_high_precision DECIMAL(38,10)
);

-- Insert test data
INSERT INTO dbo.money_test (id, description, money_col, smallmoney_col, decimal_col, decimal_high_precision) VALUES
(1, 'Small positive', 4.00, 4.00, 4.0000, 4.0000000000),
(2, 'Problematic value', 4491.00, 4491.00, 4491.0000, 4491.0000000000),
(3, 'Large value', 99999.99, 99999.99, 99999.9900, 99999.9900000000),
(4, 'Negative value', -1234.56, -1234.56, -1234.5600, -1234.5600000000),
(5, 'Zero', 0.00, 0.00, 0.0000, 0.0000000000),
(6, 'Small decimal', 0.01, 0.01, 0.0100, 0.0100000000),
(7, 'NULL values', NULL, NULL, NULL, NULL),
(8, 'Max smallmoney', 214748.3647, 214748.3647, 214748.3647, 214748.3647000000),
(9, 'Large money', 922337203685.4775, NULL, 922337203685.4775, 922337203685.4775000000),
(10, 'Precision test', 123.4567, 123.4567, 123.4567, 123.4567890123);

end:
- type: check
check: runs.dbo_money_test.status == "success"
on_failure: break

# Get source aggregations
- type: query
connection: '{source.name}'
query: |
SELECT
COUNT(*) as row_count,
SUM(CAST(money_col AS DECIMAL(38,4))) as money_sum,
SUM(CAST(smallmoney_col AS DECIMAL(38,4))) as smallmoney_sum,
SUM(decimal_col) as decimal_sum,
MIN(CAST(money_col AS DECIMAL(38,4))) as money_min,
MAX(CAST(money_col AS DECIMAL(38,4))) as money_max
FROM dbo.money_test
into: source_aggregations

# Get target aggregations
- type: query
connection: '{target.name}'
query: |
SELECT
COUNT(*) as row_count,
SUM(money_col) as money_sum,
SUM(smallmoney_col) as smallmoney_sum,
SUM(decimal_col) as decimal_sum,
MIN(money_col) as money_min,
MAX(money_col) as money_max
FROM public.money_test
into: target_aggregations

- type: log
message: |

Source (MSSQL) aggregations:
{pretty_table(store.source_aggregations)}

Target (StarRocks) aggregations:
{pretty_table(store.target_aggregations)}

# Check row count
- type: check
check: int_parse(store.source_aggregations[0].row_count) == int_parse(store.target_aggregations[0].row_count)
failure_message: Row count mismatch. Source {store.source_aggregations[0].row_count} != Target {store.target_aggregations[0].row_count}

# Check money_col sum
- type: check
check: float_parse(store.source_aggregations[0].money_sum) == float_parse(store.target_aggregations[0].money_sum)
failure_message: Money column sum mismatch. Source {store.source_aggregations[0].money_sum} != Target {store.target_aggregations[0].money_sum}

# Check decimal_col sum
- type: check
check: float_parse(store.source_aggregations[0].decimal_sum) == float_parse(store.target_aggregations[0].decimal_sum)
failure_message: Decimal column sum mismatch. Source {store.source_aggregations[0].decimal_sum} != Target {store.target_aggregations[0].decimal_sum}

# Get individual rows to check specific values
- type: query
connection: '{target.name}'
query: |
SELECT
id,
description,
money_col,
smallmoney_col,
decimal_col,
decimal_high_precision
FROM public.money_test
WHERE id IN (2, 9)
ORDER BY id
into: target_specific_rows

- type: log
message: |
Target specific rows (checking problematic value 4491.00):
{pretty_table(store.target_specific_rows)}

# Check the problematic value 4491.00
- type: check
check: float_parse(store.target_specific_rows[0].money_col) == 4491.00
failure_message: Problematic value 4491.00 not preserved correctly. Got {store.target_specific_rows[0].money_col}

# Check large money value
- type: check
check: float_parse(store.target_specific_rows[1].money_col) == 922337203685.4775
failure_message: Large money value not preserved correctly. Got {store.target_specific_rows[1].money_col}

# Get source rows for comparison
- type: query
connection: '{source.name}'
query: |
SELECT
id,
description,
CAST(money_col AS DECIMAL(38,4)) as money_value,
CAST(decimal_col AS DECIMAL(38,4)) as decimal_value
FROM dbo.money_test
ORDER BY id
into: source_rows

# Get target rows for comparison
- type: query
connection: '{target.name}'
query: |
SELECT
id,
description,
money_col as money_value,
decimal_col as decimal_value
FROM public.money_test
ORDER BY id
into: target_rows
- type: query
connection: '{target.name}'
query: |
SELECT
id,
description,
money_col as money_value,
decimal_col as decimal_value
FROM public.money_test_sql
ORDER BY id
into: target_rows_sql

- type: log
message: |
Row-by-row comparison:

Source (MSSQL) rows:
{pretty_table(store.source_rows)}

Target (StarRocks) rows:
{pretty_table(store.target_rows)}

Target (SQL) (StarRocks) rows:
{pretty_table(store.target_rows_sql)}

# Check specific rows match
- type: check
check: float_parse(store.source_rows[0].money_value) == float_parse(store.target_rows[0].money_value)
failure_message: Row 1 money value mismatch. Source {store.source_rows[0].money_value} != Target {store.target_rows[0].money_value}

- type: check
check: float_parse(store.source_rows[1].money_value) == float_parse(store.target_rows[1].money_value)
failure_message: Row 2 (4491.00) money value mismatch. Source {store.source_rows[1].money_value} != Target {store.target_rows[1].money_value}

- type: check
check: float_parse(store.source_rows[1].money_value) == float_parse(store.target_rows_sql[1].money_value)
failure_message: Row 2 (4491.00) money value mismatch. Source {store.source_rows[1].money_value} != Target {store.target_rows_sql[1].money_value}

- type: check
check: float_parse(store.source_rows[8].money_value) == float_parse(store.target_rows[8].money_value)
failure_message: Row 9 (large value) money value mismatch. Source {store.source_rows[8].money_value} != Target {store.target_rows[8].money_value}

- type: check
check: float_parse(store.source_rows[8].decimal_value) == float_parse(store.target_rows_sql[8].decimal_value)
failure_message: Row 9 (large value) money value mismatch. Source {store.source_rows[8].decimal_value} != Target {store.target_rows_sql[8].decimal_value}

# Cleanup
- type: query
connection: '{source.name}'
query: DROP TABLE IF EXISTS dbo.money_test

- type: query
connection: '{target.name}'
query: |
DROP TABLE IF EXISTS public.money_test;
DROP TABLE IF EXISTS public.money_test_sql;

streams:
dbo.money_test:
object: public.money_test
mode: full-refresh

dbo.money_test_sql:
sql: |
select * from dbo.money_test
object: public.money_test_sql
mode: full-refresh

env:
SLING_SAMPLE_SIZE: 2
Loading
Loading