构建本地模拟环境
hudi代码走读
创建本地模拟环境,一步一步调试
主程序
1 | String sourceTable = "source_01"; |
创建数据源
1 | public static String sourceDdl() { |
创建hudi表
1 | public static String sinkTableDdl(String targetTable, String basePath, String dbName) { |
go on,跟进调试
// TODO 整理到 flink demo
Flink-HoodieTableFactory


org.apache.hudi.table.HoodieTableSink
1 | // bootstrap |

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

经历三步: bootstrap → hoodieStreamWrite → compact / clean
bootstrap
- row_data_to_hoodie_record过程如下:
Flink rowData 转换为 HoodieRecord: org.apache.hudi.sink.transform.RowDataToHoodieFunction#toHoodieRecord
数据最终类型为: HoodieAvroRecord
- 如果配置了 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)

1 | org.apache.hudi.sink.transform.RowDataToHoodieFunction |
1 | // org.apache.hudi.sink.partitioner.BucketAssignFunction#processRecord |
// stream_write
// org.apache.hudi.sink.StreamWriteFunction
org.apache.hudi.sink.StreamWriteFunction#bufferRecord