hudi-agent commented on code in PR #19540:
URL: https://github.com/apache/hudi/pull/19540#discussion_r3735462212
##########
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestPipelines.java:
##########
@@ -69,6 +79,210 @@ void testGlobalRLIShufflesBucketAssignByGlobalRecordIndex()
throws Exception {
assertEquals(2, countCustomPartitions(pipeline,
GlobalRecordIndexPartitioner.class));
}
+ @Test
+ void testBootstrapPipelineSelection() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.INDEX_GLOBAL_ENABLED, false);
+ DataStream<RowData> input = rowDataInput();
+
+ DataStream<HoodieFlinkInternalRow> overwrite =
+ Pipelines.bootstrap(conf, TestConfigurations.ROW_TYPE, input, false,
true);
+ assertEquals("row_data_to_hoodie_record",
overwrite.getTransformation().getName());
+
+ DataStream<HoodieFlinkInternalRow> bounded =
+ Pipelines.bootstrap(conf, TestConfigurations.ROW_TYPE, input, true,
false);
+ assertEquals("batch_index_bootstrap",
bounded.getTransformation().getName());
+
+ conf.set(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, true);
+ conf.set(FlinkOptions.INDEX_BOOTSTRAP_TASKS, 3);
+ DataStream<HoodieFlinkInternalRow> streaming =
+ Pipelines.bootstrap(conf, TestConfigurations.ROW_TYPE, input, false,
false);
+ assertEquals("index_bootstrap", streaming.getTransformation().getName());
+ assertEquals(3, streaming.getParallelism());
+ }
+
+ @Test
+ void testWritePipelineOperatorGraphs() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.BUCKET_ASSIGN_TASKS, 2);
+ conf.set(FlinkOptions.WRITE_TASKS, 3);
+ DataStream<HoodieFlinkInternalRow> input = hoodieRowInput();
+
+ DataStream<RowData> pipeline =
+ Pipelines.hoodieStreamWrite(conf, TestConfigurations.ROW_TYPE, input);
+
+ assertEquals(3, pipeline.getParallelism());
+ assertTrue(transformationNames(pipeline).stream().anyMatch(name ->
name.startsWith("bucket_assigner")));
+ assertTrue(transformationNames(pipeline).stream().anyMatch(name ->
name.startsWith("stream_write:")));
+ }
+
+ @Test
+ void testSimpleBucketWritePipeline() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+ conf.set(FlinkOptions.BUCKET_INDEX_ENGINE_TYPE,
+ HoodieIndex.BucketIndexEngineType.SIMPLE.name());
+ conf.set(FlinkOptions.WRITE_TASKS, 4);
+
+ DataStream<RowData> pipeline = Pipelines.hoodieStreamWrite(
+ conf, TestConfigurations.ROW_TYPE, hoodieRowInput());
+
+ assertEquals(4, pipeline.getParallelism());
+
assertTrue(pipeline.getTransformation().getName().startsWith("bucket_write:"));
+ }
+
+ @Test
+ void testBulkInsertAndAppendValidation() {
+ Configuration recordIndexConf = defaultConf();
+ recordIndexConf.set(FlinkOptions.INDEX_TYPE,
HoodieIndex.IndexType.RECORD_LEVEL_INDEX.name());
+ assertThrows(HoodieException.class,
+ () -> Pipelines.bulkInsert(recordIndexConf,
TestConfigurations.ROW_TYPE, rowDataInput()));
+
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+ conf.set(FlinkOptions.BUCKET_INDEX_ENGINE_TYPE,
+ HoodieIndex.BucketIndexEngineType.CONSISTENT_HASHING.name());
+ Configuration consistentBucketConf = conf;
+ assertThrows(HoodieException.class,
+ () -> Pipelines.bulkInsert(
+ consistentBucketConf, TestConfigurations.ROW_TYPE,
rowDataInput()));
+ assertThrows(HoodieNotSupportedException.class,
+ () -> Pipelines.append(
+ consistentBucketConf, TestConfigurations.ROW_TYPE,
rowDataInput()));
+ }
+
+ @Test
+ void testBulkInsertGraphForDefaultAndLsmLayouts() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.OPERATION, WriteOperationType.BULK_INSERT.value());
+ conf.set(FlinkOptions.WRITE_TASKS, 2);
+ DataStream<RowData> lsmPipeline =
+ Pipelines.bulkInsert(conf, TestConfigurations.ROW_TYPE,
rowDataInput());
+ assertEquals(2, lsmPipeline.getParallelism());
+
assertTrue(transformationNames(lsmPipeline).contains("lsm_bulk_insert_sort_keys"));
+
assertTrue(transformationNames(lsmPipeline).contains("lsm_sorter:(partition_path,
record_key)"));
+
+ conf = defaultConf();
+ conf.set(FlinkOptions.OPERATION, WriteOperationType.INSERT.value());
+ conf.set(FlinkOptions.WRITE_TASKS, 3);
+ conf.set(FlinkOptions.WRITE_BULK_INSERT_SORT_INPUT, true);
+ DataStream<RowData> defaultPipeline =
+ Pipelines.bulkInsert(conf, TestConfigurations.ROW_TYPE,
rowDataInput());
+ assertEquals(3, defaultPipeline.getParallelism());
+
assertTrue(transformationNames(defaultPipeline).contains("sorter:(partition_key)"));
+ }
+
+ @Test
+ void testBucketBulkInsertGraphs() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+ conf.set(FlinkOptions.BUCKET_INDEX_ENGINE_TYPE,
+ HoodieIndex.BucketIndexEngineType.SIMPLE.name());
+ conf.set(FlinkOptions.OPERATION, WriteOperationType.INSERT.value());
+ conf.set(FlinkOptions.WRITE_TASKS, 2);
+ conf.set(FlinkOptions.WRITE_BULK_INSERT_SORT_INPUT, false);
+
+ DataStream<RowData> unsorted =
+ Pipelines.bulkInsert(conf, TestConfigurations.ROW_TYPE,
rowDataInput());
+ assertEquals(2, unsorted.getParallelism());
+
assertTrue(unsorted.getTransformation().getName().startsWith("bucket_bulk_insert"));
+
+ conf.set(FlinkOptions.WRITE_BULK_INSERT_SORT_INPUT, true);
+ DataStream<RowData> sorted =
+ Pipelines.bulkInsert(conf, TestConfigurations.ROW_TYPE,
rowDataInput());
+ assertTrue(transformationNames(sorted).contains("file_sorter"));
+
+ Configuration lsmConf = defaultConf();
+ lsmConf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+ lsmConf.set(FlinkOptions.BUCKET_INDEX_ENGINE_TYPE,
+ HoodieIndex.BucketIndexEngineType.SIMPLE.name());
+ lsmConf.set(FlinkOptions.OPERATION,
WriteOperationType.BULK_INSERT.value());
+ lsmConf.set(FlinkOptions.WRITE_TASKS, 3);
+
+ DataStream<RowData> lsm =
+ Pipelines.bulkInsert(lsmConf, TestConfigurations.ROW_TYPE,
rowDataInput());
+ assertEquals(3, lsm.getParallelism());
+ assertTrue(transformationNames(lsm).contains("lsm_bulk_insert_sort_keys"));
+ assertTrue(transformationNames(lsm).contains("lsm_sorter:(file_group,
record_key)"));
+ }
+
+ @Test
+ void testServiceAndSinkGraphs() {
+ Configuration conf = defaultConf();
+ conf.set(FlinkOptions.COMPACTION_TASKS, 4);
+ conf.set(FlinkOptions.CLUSTERING_TASKS, 3);
+ DataStream<RowData> input = rowDataInput();
+
+ DataStreamSink<?> compaction = Pipelines.compact(conf, input);
+ assertEquals("compact_commit", compaction.getTransformation().getName());
+ assertEquals(1, compaction.getTransformation().getParallelism());
+ assertEquals(1, compaction.getTransformation().getMaxParallelism());
+
+ DataStreamSink<?> clustering = Pipelines.cluster(conf,
TestConfigurations.ROW_TYPE, input);
+ assertEquals("clustering_commit",
clustering.getTransformation().getName());
+ assertEquals(1, clustering.getTransformation().getParallelism());
+ assertEquals(1, clustering.getTransformation().getMaxParallelism());
+
+ DataStreamSink<?> clean = Pipelines.clean(conf, input);
+ assertEquals("clean_commits", clean.getTransformation().getName());
+ assertEquals(1, clean.getTransformation().getParallelism());
+ assertEquals(1, clean.getTransformation().getMaxParallelism());
+
+ DataStreamSink<?> dummy = Pipelines.dummySink(rowDataInput(5));
+ assertEquals("dummy", dummy.getTransformation().getName());
+ assertEquals(5, dummy.getTransformation().getParallelism());
+ }
+
+ @Test
+ void testOperatorNamesAndIndexPartitioner() {
+ Configuration conf = defaultConf();
+ assertEquals("write: analytics.orders", Pipelines.opName("write", conf));
+ String firstUid = Pipelines.opUID("unique_test_operator", conf);
+ String secondUid = Pipelines.opUID("unique_test_operator", conf);
+
assertTrue(firstUid.startsWith("uid_unique_test_operator_analytics.orders"));
+
assertTrue(secondUid.startsWith("uid_unique_test_operator_1_analytics.orders"));
+ assertNotEquals(firstUid, secondUid);
+ assertEquals(2, new Pipelines.IndexPartitioner().partition(8, 3));
+ }
+
+ private Configuration defaultConf() {
+ Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+ conf.set(FlinkOptions.DATABASE_NAME, "analytics");
+ conf.set(FlinkOptions.TABLE_NAME, "orders");
+ return conf;
+ }
+
+ private DataStream<RowData> rowDataInput() {
+ return rowDataInput(1);
+ }
+
+ private DataStream<RowData> rowDataInput(int parallelism) {
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
+ DataStreamSource<RowData> source = env.fromCollection(
+ Collections.<RowData>emptyList(),
+
org.apache.flink.table.runtime.typeutils.InternalTypeInfo.of(TestConfigurations.ROW_TYPE));
+ if (parallelism == 1) {
+ return source;
+ }
+ return source.map(
+ row -> row,
+
org.apache.flink.table.runtime.typeutils.InternalTypeInfo.of(TestConfigurations.ROW_TYPE))
+ .setParallelism(parallelism);
+ }
+
+ private DataStream<HoodieFlinkInternalRow> hoodieRowInput() {
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
+ return env.fromCollection(
Review Comment:
🤖 nit: the fully-qualified
`org.apache.flink.table.runtime.typeutils.InternalTypeInfo` is repeated inline
here and below — could you add an import so these calls read more cleanly?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]