Skip to content

Commit de5caea

Browse files
committed
Persist aggregate state settings for plain exports
Carry the Parquet and Iceberg aggregate-state gates through plain `MergeTree` export task serialization and worker contexts. Keep older descriptors compatible by treating missing fields as disabled.
1 parent ba43db9 commit de5caea

3 files changed

Lines changed: 79 additions & 0 deletions

File tree

‎src/Storages/MergeTree/MergeTreePartitionExportTask.h‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,8 @@ struct MergeTreePartitionExportTask
8282
String filename_pattern;
8383
bool write_full_path_in_iceberg_metadata = false;
8484
bool allow_lossy_cast = false;
85+
bool allow_aggregate_function_states_in_parquet = false;
86+
bool allow_aggregate_function_states_in_iceberg = false;
8587
String iceberg_metadata_json;
8688

8789
/// Optional for backwards compatibility with descriptors written before these
@@ -189,6 +191,8 @@ struct MergeTreePartitionExportTask
189191
json.set("filename_pattern", filename_pattern);
190192
json.set("write_full_path_in_iceberg_metadata", write_full_path_in_iceberg_metadata);
191193
json.set("allow_lossy_cast", allow_lossy_cast);
194+
json.set("allow_aggregate_function_states_in_parquet", allow_aggregate_function_states_in_parquet);
195+
json.set("allow_aggregate_function_states_in_iceberg", allow_aggregate_function_states_in_iceberg);
192196
if (!iceberg_metadata_json.empty())
193197
json.set("iceberg_metadata_json", iceberg_metadata_json);
194198
if (parquet_compression_method)
@@ -281,6 +285,10 @@ struct MergeTreePartitionExportTask
281285
task.filename_pattern = json->getValue<String>("filename_pattern");
282286
task.write_full_path_in_iceberg_metadata = json->getValue<bool>("write_full_path_in_iceberg_metadata");
283287
task.allow_lossy_cast = json->getValue<bool>("allow_lossy_cast");
288+
if (json->has("allow_aggregate_function_states_in_parquet"))
289+
task.allow_aggregate_function_states_in_parquet = json->getValue<bool>("allow_aggregate_function_states_in_parquet");
290+
if (json->has("allow_aggregate_function_states_in_iceberg"))
291+
task.allow_aggregate_function_states_in_iceberg = json->getValue<bool>("allow_aggregate_function_states_in_iceberg");
284292
if (json->has("iceberg_metadata_json"))
285293
task.iceberg_metadata_json = json->getValue<String>("iceberg_metadata_json");
286294
if (json->has("parquet_compression_method"))

‎src/Storages/MergeTree/tests/gtest_export_partition_ordering.cpp‎

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
#include <sstream>
44
#include <Storages/ExportReplicatedMergeTreePartitionTaskEntry.h>
55
#include <Storages/MergeTree/ExportPartitionUtils.h>
6+
#include <Storages/MergeTree/MergeTreePartitionExportTask.h>
67
#include <Common/tests/gtest_global_context.h>
78
#include <Core/Settings.h>
89
#include <Poco/JSON/Object.h>
@@ -51,6 +52,29 @@ namespace
5152
manifest.filename_pattern = "{part_name}";
5253
return manifest;
5354
}
55+
56+
MergeTreePartitionExportTask makeValidPlainTask()
57+
{
58+
MergeTreePartitionExportTask task;
59+
task.transaction_id = "tx1";
60+
task.query_id = "query1";
61+
task.partition_id = "2020";
62+
task.source_database = "db1";
63+
task.source_table = "source1";
64+
task.destination_database = "db1";
65+
task.destination_table = "table1";
66+
task.create_time = 1000;
67+
task.parts.push_back({"part1", false, {}});
68+
task.task_timeout_seconds = 60;
69+
task.max_threads = 1;
70+
task.parallel_formatting = true;
71+
task.parquet_parallel_encoding = true;
72+
task.max_bytes_per_file = 1000000;
73+
task.max_rows_per_file = 1000;
74+
task.file_already_exists_policy = MergeTreePartExportManifest::FileAlreadyExistsPolicy::error;
75+
task.filename_pattern = "{part_name}";
76+
return task;
77+
}
5478
}
5579

5680
class ExportPartitionOrderingTest : public ::testing::Test
@@ -270,6 +294,47 @@ TEST_F(ExportPartitionManifestBackCompatTest, AggregateFunctionStateGatesApplied
270294
}
271295
}
272296

297+
TEST_F(ExportPartitionManifestBackCompatTest, PlainTaskMissingAggregateFunctionStateGatesParseAsDisabled)
298+
{
299+
auto task = makeValidPlainTask();
300+
task.allow_aggregate_function_states_in_parquet = true;
301+
task.allow_aggregate_function_states_in_iceberg = true;
302+
303+
Poco::JSON::Parser parser;
304+
auto json = parser.parse(task.toJsonString()).extract<Poco::JSON::Object::Ptr>();
305+
json->remove("allow_aggregate_function_states_in_parquet");
306+
json->remove("allow_aggregate_function_states_in_iceberg");
307+
std::ostringstream oss; // STYLE_CHECK_ALLOW_STD_STRING_STREAM
308+
oss.exceptions(std::ios::failbit);
309+
Poco::JSON::Stringifier::stringify(json, oss);
310+
311+
auto parsed = MergeTreePartitionExportTask::fromJsonString(oss.str());
312+
EXPECT_FALSE(parsed.allow_aggregate_function_states_in_parquet);
313+
EXPECT_FALSE(parsed.allow_aggregate_function_states_in_iceberg);
314+
}
315+
316+
TEST_F(ExportPartitionManifestBackCompatTest, PlainTaskAggregateFunctionStateGatesAppliedToWorkerContextForEveryValue)
317+
{
318+
for (const bool value : {false, true})
319+
{
320+
auto task = makeValidPlainTask();
321+
task.allow_aggregate_function_states_in_parquet = value;
322+
task.allow_aggregate_function_states_in_iceberg = value;
323+
324+
auto parsed = MergeTreePartitionExportTask::fromJsonString(task.toJsonString());
325+
auto worker_context = ExportPartitionUtils::getContextCopyWithTaskSettings(getContext().context, parsed);
326+
327+
EXPECT_EQ(parsed.allow_aggregate_function_states_in_parquet, value) << "value=" << value;
328+
EXPECT_EQ(parsed.allow_aggregate_function_states_in_iceberg, value) << "value=" << value;
329+
EXPECT_EQ(
330+
worker_context->getSettingsRef()[Setting::allow_experimental_aggregate_function_states_in_parquet].value,
331+
value) << "value=" << value;
332+
EXPECT_EQ(
333+
worker_context->getSettingsRef()[Setting::allow_experimental_aggregate_function_states_in_iceberg].value,
334+
value) << "value=" << value;
335+
}
336+
}
337+
273338
TEST(ExportPartitionRetryClassification, MissingPartIsFatalOnlyOnPlainPath)
274339
{
275340
EXPECT_FALSE(ExportPartitionUtils::isNonRetryableExportError(ErrorCodes::NO_SUCH_DATA_PART));

‎src/Storages/StorageMergeTree.cpp‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -138,6 +138,8 @@ namespace Setting
138138
extern const SettingsTimezone iceberg_partition_timezone;
139139
extern const SettingsMergeTreePartExportSchemaMatchMode export_merge_tree_part_schema_match_mode;
140140
extern const SettingsBool export_merge_tree_part_ignore_extra_source_columns;
141+
extern const SettingsBool allow_experimental_aggregate_function_states_in_parquet;
142+
extern const SettingsBool allow_experimental_aggregate_function_states_in_iceberg;
141143
}
142144

143145
namespace ServerSetting
@@ -3894,6 +3896,10 @@ void StorageMergeTree::exportPartitionToTable(const PartitionCommand & command,
38943896
descriptor.filename_pattern = query_context->getSettingsRef()[Setting::export_merge_tree_part_filename_pattern].value;
38953897
descriptor.write_full_path_in_iceberg_metadata = query_context->getSettingsRef()[Setting::write_full_path_in_iceberg_metadata];
38963898
descriptor.allow_lossy_cast = query_context->getSettingsRef()[Setting::export_merge_tree_part_allow_lossy_cast];
3899+
descriptor.allow_aggregate_function_states_in_parquet
3900+
= query_context->getSettingsRef()[Setting::allow_experimental_aggregate_function_states_in_parquet];
3901+
descriptor.allow_aggregate_function_states_in_iceberg
3902+
= query_context->getSettingsRef()[Setting::allow_experimental_aggregate_function_states_in_iceberg];
38973903
descriptor.iceberg_partition_timezone = query_context->getSettingsRef()[Setting::iceberg_partition_timezone].toString();
38983904
descriptor.schema_match_mode = query_context->getSettingsRef()[Setting::export_merge_tree_part_schema_match_mode].value;
38993905
descriptor.ignore_extra_source_columns = query_context->getSettingsRef()[Setting::export_merge_tree_part_ignore_extra_source_columns].value;

0 commit comments

Comments
 (0)