0%

hudi源码走读-Flink数据写入

构建本地模拟环境

hudi代码走读

创建本地模拟环境,一步一步调试

主程序

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
String sourceTable = "source_01";
ParameterTool pt = ParameterTool.fromArgs(args);
String database = pt.get("database", "ods_glory");
String table = pt.get("table", "ods_glory_t1");
String hudiDemoHome = "apps/data/hudi";
String tablePath = String.format("%s/%s/%s.db/%s", System.getProperty("user.home"), hudiDemoHome, database, table);
log.info("userParameter: {}", pt.toMap());
//
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
env.setParallelism(1);
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);
Configuration configuration = tableEnv.getConfig().getConfiguration();
configuration.setString("table.dynamic-table-options.enabled", "true");
configuration.setString(PipelineOptions.NAME.key(), "flake-hudi-1.13");
log.info("envConfig: {}.", configuration.toMap());
//
tableEnv.executeSql(sourceDdl());
String targetDal = sinkTableDdl(table, tablePath, database);
log.info("targetDal:{}.", targetDal);
tableEnv.executeSql(targetDal);
// tableEnv.executeSql(String.format("select *,DATE_FORMAT(ts, 'yyyyMMdd') as dt from %s", table, sourceTable))
// System.out.println(tableEnv.explain(tableEnv.from(table)));
String sql = String.format("insert into %s select *,DATE_FORMAT(ts, 'yyyyMMdd') as dt from %s", table,
sourceTable);
tableEnv.executeSql(sql);
System.out.println(env.getExecutionPlan());

创建数据源

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public static String sourceDdl() {
String sourceTable = "source_01";
return String.format("CREATE TABLE %s (\n" +
" uuid varchar(20),\n" +
" facc_time STRING ,\n" +
" facc_time_rongduan STRING,\n" +
" facc_type BIGINT ,\n" +
" fact_info STRING ,\n" +
" ts timestamp(3)\n" +
") WITH (\n" +
" 'connector' = 'datagen',\n" +
" 'rows-per-second' = '100'\n" +
")", sourceTable);
}

创建hudi表

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
public static String sinkTableDdl(String targetTable, String basePath, String dbName) {
return String.format("create table %s(\n" +
" uuid STRING,\n" +
" facc_time STRING ,\n" +
" facc_time_rongduan STRING,\n" +
" facc_type BIGINT ,\n" +
" fact_info STRING ,\n" +
" ts timestamp(3),\n" +
" dt string,\n" +
" PRIMARY KEY(uuid) NOT ENfORCED" +
")\n" +
" PARTITIONED BY (`dt`)\n" +
" with (\n" +
" 'connector' = 'hudi',\n" +
" 'path' = '%s', -- 替换成的绝对路径\n" +
" 'table.type' = 'MERGE_ON_READ',\n" +
" 'write.bucket_assign.tasks' = '8',\n" +
" 'write.tasks' = '8',\n" +
" 'write.operation' = 'upsert', -- upsert/insert\n" +
" 'changelog.enabled' = 'true',\n" +
" 'read.streaming.enabled' = 'true',\n" +
" 'read.streaming.check-interval' = '1',\n" +
" 'compaction.tasks' = '8',\n" +
" 'compaction.trigger.strategy'='num_commits',\n" +
" 'compaction.delta_commits' ='5',\n" +
" 'compaction.max_memory' = '3096',\n" +
" 'clean.retain_commits' = '30',\n" +
" 'hive_sync.enable' = 'false',\n" +
" 'hive_sync.mode' = 'hms',\n" +
" 'hive_sync.db' = '%s',\n" +
" 'hive_sync.table' = '%s',\n" +
" 'hive_sync.metastore.uris' = '%s'\n" +
")", targetTable, basePath, dbName, targetTable, "metastoreUrl");
}

go on,跟进调试

// TODO 整理到 flink demo

img

img

org.apache.hudi.table.HoodieTableSink

1
2
3
4
5
6
7
8
9
10
11
// bootstrap
final DataStream<HoodieRecord> hoodieRecordDataStream =
Pipelines.bootstrap(conf, rowType, parallelism, dataStream, context.isBounded(), overwrite);
// write pipeline
DataStream<Object> pipeline = Pipelines.hoodieStreamWrite(conf, parallelism, hoodieRecordDataStream);
// compaction
if (StreamerUtil.needsAsyncCompaction(conf)) {
return Pipelines.compact(conf, pipeline);
} else {
return Pipelines.clean(conf, pipeline);
}

img

org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecSink#createSinkTransformation

img

经历三步: bootstrap → hoodieStreamWrite → compact / clean

bootstrap

  1. row_data_to_hoodie_record过程如下:

Flink rowData 转换为 HoodieRecord: org.apache.hudi.sink.transform.RowDataToHoodieFunction#toHoodieRecord

数据最终类型为: HoodieAvroRecord

  1. 如果配置了 index.bootstrap.enabled,会增加一个 index_bootstrap 节点,用于在flink state中保存instant信息(lastInstantTime)

StreamWrite

org.apache.hudi.sink.utils.Pipelines#hoodieStreamWrite

compact

clean

org.apache.hudi.sink.utils.Pipelines#bootstrap(org.apache.flink.configuration.Configuration, org.apache.flink.table.types.logical.RowType, int, org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.data.RowData>, boolean, boolean)

img

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
org.apache.hudi.sink.transform.RowDataToHoodieFunction

org.apache.hudi.sink.bootstrap.BootstrapOperator 这个需要继续分析

// org.apache.hudi.sink.utils.Pipelines#hoodieStreamWrite
dataStream
.keyBy(HoodieRecord::getRecordKey) // hoodieKey 需要单独分析
.transform("bucket_assigner", TypeInformation.of(HoodieRecord.class), new KeyedProcessOperator(new BucketAssignFunction(conf)))
// org.apache.hudi.sink.partitioner.BucketAssignFunction#processRecord
// hoodieKey 的 recordKey 和 partitionPath 代表是什么意思?
// org.apache.hudi.common.model.HoodieRecordLocation 和 org.apache.hudi.common.model.HoodieRecordGlobalLocation
.uid("uid_bucket_assigner_" + conf.getString(FlinkOptions.TABLE_NAME))
.setParallelism((Integer)conf.getOptional(FlinkOptions.BUCKET_ASSIGN_TASKS).orElse(defaultParallelism))
.keyBy((record) -> { return record.getCurrentLocation().getFileId(); })
.transform("stream_write", TypeInformation.of(Object.class), operatorFactory)
// org.apache.flink.streaming.api.operators.SimpleUdfStreamOperatorFactory
// Flink的单UDF的StreamOperatorFactory
// org.apache.hudi.sink.common.WriteOperatorFactory#createStreamOperator
.uid("uid_stream_write" + conf.getString(FlinkOptions.TABLE_NAME))
.setParallelism(conf.getInteger(FlinkOptions.WRITE_TASKS));

dataStream
.transform("compact_plan_generate", TypeInformation.of(CompactionPlanEvent.class), new CompactionPlanOperator(conf))
.setParallelism(1)
.rebalance()
.transform("compact_task", TypeInformation.of(CompactionCommitEvent.class), new ProcessOperator(new CompactFunction(conf)))
.setParallelism(conf.getInteger(FlinkOptions.COMPACTION_TASKS))
.addSink(new CompactionCommitSink(conf))
.name("compact_commit")
.setParallelism(1);
1
2
3
4
5
// org.apache.hudi.sink.partitioner.BucketAssignFunction#processRecord
HoodieRecord<?> deleteRecord = new HoodieAvroRecord(new HoodieKey(recordKey, oldLoc.getPartitionPath()), this.payloadCreation.createDeletePayload((BaseAvroPayload)record.getData()));
deleteRecord.setCurrentLocation(oldLoc.toLocal("U"));
deleteRecord.seal();
out.collect(deleteRecord);

// stream_write
// org.apache.hudi.sink.StreamWriteFunction
org.apache.hudi.sink.StreamWriteFunction#bufferRecord