diff options
| author | artemmashin <[email protected]> | 2026-07-16 14:53:41 +0300 |
|---|---|---|
| committer | artemmashin <[email protected]> | 2026-07-17 12:58:08 +0300 |
| commit | 3dff175fc5960889dfbf66bdbd9f73e31d784068 (patch) | |
| tree | 59a82d4ea72988cc37dd885024f148093b76a192 | |
| parent | 6803bcb87218573d48e07ee929f7c93bf7d8b3ac (diff) | |
Support proper write to sorted yt table
commit_hash:572af6eec0409e60c1c7a234f7b60275da0aaf37
18 files changed, 190 insertions, 0 deletions
diff --git a/yt/yql/providers/yt/provider/yql_yt_ytflow_integration.cpp b/yt/yql/providers/yt/provider/yql_yt_ytflow_integration.cpp index 9e4c87017fe..c733740bf06 100644 --- a/yt/yql/providers/yt/provider/yql_yt_ytflow_integration.cpp +++ b/yt/yql/providers/yt/provider/yql_yt_ytflow_integration.cpp @@ -289,6 +289,15 @@ public: sinkSettings.SetDoesExist(tableDesc.Meta->DoesExist); sinkSettings.SetTruncate(tableDesc.Intents & TYtTableIntent::Override); + + const auto& rowSpec = tableDesc.RowSpec; + if (rowSpec) { + for (const auto& [column, _] : rowSpec->GetForeignSort()) { + if (!rowSpec->ExpressionColumns.contains(column)) { + sinkSettings.AddKeyColumns(column); + } + } + } } settings.PackFrom(sinkSettings); diff --git a/yt/yql/providers/ytflow/integration/proto/yt.proto b/yt/yql/providers/ytflow/integration/proto/yt.proto index 2959ca4b6c8..4ebfad1059b 100644 --- a/yt/yql/providers/ytflow/integration/proto/yt.proto +++ b/yt/yql/providers/ytflow/integration/proto/yt.proto @@ -14,4 +14,5 @@ message TQYTSinkMessage { optional bytes RowType = 3; optional bool DoesExist = 5; optional bool Truncate = 6; + repeated string KeyColumns = 7; } diff --git a/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.cfg b/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.cfg new file mode 100644 index 00000000000..0e3721e4b77 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.cfg @@ -0,0 +1,4 @@ +in Input input.txt +out OutputKey output_sorted_by_key.txt +out OutputValue output_sorted_by_value.txt +providers ytflow diff --git a/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.yql b/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.yql new file mode 100644 index 00000000000..b938bbf058a --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/multiple_sorted_outputs_with_different_keys.yql @@ -0,0 +1,27 @@ +/* sort outputs */ + +use plato; + +pragma Engine = "ytflow"; + +pragma Ytflow.Cluster = "plato"; +pragma Ytflow.PipelinePath = "pipelines/test"; + +$lambda = ($row) -> { + $output_row_type = Struct<"key":Optional<String>, "value":Optional<Int64>>; + $variant_type = Variant<$output_row_type, $output_row_type>; + + return If( + $row.int64_field < 10, + Variant(<|"key":$row.string_field, "value":$row.int64_field|>, "0", $variant_type), + Variant(<|"key":$row.string_field, "value":$row.int64_field|>, "1", $variant_type) + ); +}; + +$sorted_by_key, $sorted_by_value = process Input using $lambda(TableRow()); + +insert into OutputKey +select * from $sorted_by_key; + +insert into OutputValue +select * from $sorted_by_value; diff --git a/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt b/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt new file mode 100644 index 00000000000..e69de29bb2d --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt diff --git a/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt.attr b/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt.attr new file mode 100644 index 00000000000..8f69a0c4304 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_composite_key_value_table.txt.attr @@ -0,0 +1,23 @@ +{ + "schema" = < + unique_keys = %true; + strict = %true; + > [ + { + "name" = "key_b"; + "type" = "string"; + "sort_order" = "ascending"; + }; + { + "name" = "key_a"; + "type" = "int64"; + "sort_order" = "ascending"; + }; + { + "name" = "kv_value"; + "type" = "int64"; + }; + ]; + "_yql_dynamic" = %true; + "_yql_dynamic_native_read" = %true; +} diff --git a/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt new file mode 100644 index 00000000000..e69de29bb2d --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt diff --git a/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt.attr b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt.attr new file mode 100644 index 00000000000..6086a84dbc9 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_key.txt.attr @@ -0,0 +1,18 @@ +{ + "schema" = < + unique_keys = %true; + strict = %true; + > [ + { + "name" = "key"; + "type" = "string"; + "sort_order" = "ascending"; + }; + { + "name" = "value"; + "type" = "int64"; + }; + ]; + "_yql_dynamic" = %true; + "_yql_dynamic_native_read" = %true; +} diff --git a/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt new file mode 100644 index 00000000000..e69de29bb2d --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt diff --git a/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt.attr b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt.attr new file mode 100644 index 00000000000..3acb927d6d8 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/output_sorted_by_value.txt.attr @@ -0,0 +1,18 @@ +{ + "schema" = < + unique_keys = %true; + strict = %true; + > [ + { + "name" = "value"; + "type" = "int64"; + "sort_order" = "ascending"; + }; + { + "name" = "key"; + "type" = "string"; + }; + ]; + "_yql_dynamic" = %true; + "_yql_dynamic_native_read" = %true; +} diff --git a/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.cfg b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.cfg new file mode 100644 index 00000000000..0e3721e4b77 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.cfg @@ -0,0 +1,4 @@ +in Input input.txt +out OutputKey output_sorted_by_key.txt +out OutputValue output_sorted_by_value.txt +providers ytflow diff --git a/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.yql b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.yql new file mode 100644 index 00000000000..7131563360a --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_different_keys.yql @@ -0,0 +1,16 @@ +/* sort outputs */ + +use plato; + +pragma Engine = "ytflow"; + +pragma Ytflow.Cluster = "plato"; +pragma Ytflow.PipelinePath = "pipelines/test"; + +$result_stream = select string_field as key, int64_field as value from Input; + +insert into OutputKey +select * from $result_stream; + +insert into OutputValue +select * from $result_stream; diff --git a/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.cfg b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.cfg new file mode 100644 index 00000000000..596696392fb --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.cfg @@ -0,0 +1,4 @@ +in Input input.txt +out OutputFirst output_sorted_by_key.txt +out OutputSecond output_sorted_by_key.txt +providers ytflow diff --git a/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.yql b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.yql new file mode 100644 index 00000000000..784e88094d3 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/reuse_output_to_sorted_tables_with_same_key.yql @@ -0,0 +1,16 @@ +/* sort outputs */ + +use plato; + +pragma Engine = "ytflow"; + +pragma Ytflow.Cluster = "plato"; +pragma Ytflow.PipelinePath = "pipelines/test"; + +$result_stream = select string_field as key, int64_field as value from Input; + +insert into OutputFirst +select * from $result_stream; + +insert into OutputSecond +select * from $result_stream; diff --git a/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.cfg b/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.cfg new file mode 100644 index 00000000000..d28699e18b4 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.cfg @@ -0,0 +1,4 @@ +in Input input.txt +out OutputSorted output_sorted_by_key.txt +out OutputUnsorted output_value_multiple_outputs.txt +providers ytflow diff --git a/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.yql b/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.yql new file mode 100644 index 00000000000..b0ab691b823 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/sorted_and_unsorted_outputs.yql @@ -0,0 +1,28 @@ +/* sort outputs */ + +use plato; + +pragma Engine = "ytflow"; + +pragma Ytflow.Cluster = "plato"; +pragma Ytflow.PipelinePath = "pipelines/test"; + +$lambda = ($row) -> { + $sorted_row_type = Struct<"key":Optional<String>, "value":Optional<Int64>>; + $unsorted_row_type = Struct<"value":Optional<Int64>>; + $variant_type = Variant<$sorted_row_type, $unsorted_row_type>; + + return If( + $row.int64_field < 10, + Variant(<|"key":$row.string_field, "value":$row.int64_field|>, "0", $variant_type), + Variant(<|"value":$row.int64_field|>, "1", $variant_type) + ); +}; + +$sorted, $unsorted = process Input using $lambda(TableRow()); + +insert into OutputSorted +select * from $sorted; + +insert into OutputUnsorted +select * from $unsorted; diff --git a/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.cfg b/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.cfg new file mode 100644 index 00000000000..ae0ae973f2e --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.cfg @@ -0,0 +1,3 @@ +in Input input_stream_composite_key.txt +out Output output_composite_key_value_table.txt +providers ytflow diff --git a/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.yql b/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.yql new file mode 100644 index 00000000000..37f5aa55876 --- /dev/null +++ b/yt/yql/tests/sql/suites/ytflow/write_to_sorted_table_with_composite_key.yql @@ -0,0 +1,15 @@ +/* sort outputs */ + +use plato; + +pragma Engine = "ytflow"; + +pragma Ytflow.Cluster = "plato"; +pragma Ytflow.PipelinePath = "pipelines/test"; + +insert into Output +select + key_a, + key_b || "_b" as key_b, + value as kv_value +from Input; |
