尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

功能点 9:Flink Connector

功能点 9:Flink Connector 功能点 9Flink Connector —— 源码阅读笔记对应源码阅读计划功能点 9FlussCatalog、FlussSource/SourceEnumerator/SourceReader、FlussSink/Writer/Committer、LookupFunction。笔记 9.1FlussCatalog —— Flink Catalog 实现文件FlussCatalog.java、FlussCatalogFactory.java路径fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/catalog/Catalog 注册入口publicclassFlussCatalogFactoryimplementsCatalogFactory{OverridepublicCatalogcreateCatalog(Stringname,MapString,Stringoptions){// 用户 SQL: CREATE CATALOG fluss WITH (typefluss, bootstrap.servers...)// Flink 调用此 Factory 创建 Catalog 实例returnnewFlussCatalog(name,options.get(bootstrap.servers),FlinkConnectorOptions.fromMap(options));}}getTable 核心实现publicclassFlussCatalogextendsAbstractCatalog{OverridepublicCatalogBaseTablegetTable(ObjectPathtablePath){// 1. ★ 通过 RPC 从 Fluss Coordinator 获取 TableDescriptorTableDescriptorflussDescadminClient.getTable(tablePath);// 2. ★ Fluss Schema → Flink Schema 转换SchemaflinkSchemaSchema.newBuilder().fromRowDataType(toFlinkDataType(flussDesc.getSchema())).build();// 3. 根据表类型创建不同的 Connector TableMapString,StringtableOptionsbuildTableOptions(flussDesc);if(flussDesc.getTableType()TableType.PRIMARY_KEY){// PK 表支持 Changelog Mode (UPSERT)returnCatalogTable.of(flinkSchema,Fluss Primary Key Table: tablePath,flussDesc.getPartitionKeys(),tableOptions);}else{// Log 表仅追加returnCatalogTable.of(flinkSchema,Fluss Log Table: tablePath,flussDesc.getPartitionKeys(),tableOptions);}}/** * ★ 关键Fluss 数据类型 → Flink 数据类型映射 */privateDataTypetoFlinkDataType(SchemaflussSchema){// Fluss Type → Flink Type// ─────────────────────────────────// INT → DataTypes.INT()// BIGINT → DataTypes.BIGINT()// STRING → DataTypes.STRING()// DECIMAL(p,s) → DataTypes.DECIMAL(p,s)// TIMESTAMP → DataTypes.TIMESTAMP(3)// ARRAYT → DataTypes.ARRAY(toFlinkType(T))RowTyperowTypeflussSchema.toRowType();// ...}}笔记 9.2FlussSource —— Source 实现文件FlussSource.java、FlussSourceEnumerator.java、FlussSourceReader.javaSource Split 定义/** * Fluss Source Split TablePath PartitionId BucketId StartOffset */publicclassFlussSourceSplitimplementsSourceSplit{privatefinalTablePathtablePath;privatefinallongpartitionId;privatefinalintbucketId;privatefinallongstartOffset;// 从这个 offset 开始读privatefinallongstopOffset;// 读到这个 offset可选-1 表示无限}Split 发现EnumeratorpublicclassFlussSourceEnumeratorimplementsSplitEnumeratorFlussSourceSplit{Overridepublicvoidstart(){// 1. 从 Coordinator 获取所有 PartitionListPartitionInfopartitionscoordinatorClient.listPartitions(sourceTable);// 2. 对每个 Partition获取所有 Bucket 信息for(PartitionInfopartition:partitions){for(intbucketId0;bucketIdpartition.getBucketCount();bucketId){longleaderServerpartition.getBucketLeader(bucketId);// 3. 创建 SplitFlussSourceSplitsplitnewFlussSourceSplit(sourceTable,partition.getPartitionId(),bucketId,discoverStartOffset(partition,bucketId)// 从 Checkpoint 恢复);pendingSplits.add(split);}}// 4. ★ 批量分配 Split 给 Reader避免逐个分配的开销assignSplitsInBatches();}/** * ★ 本地优先分配策略 * 优先将 Split 分配给与 TabletServer 在同一节点的 Reader */privatevoidassignSplitsInBatches(){MapString,ListFlussSourceSplitreaderAssignmentsnewHashMap();for(FlussSourceSplitsplit:pendingSplits){// 获取 Split 的 Leader TabletServer 地址StringleaderHostgetLeaderHost(split);// 优先分配给同一主机的 ReaderStringpreferredReaderfindLocalReader(leaderHost);readerAssignments.computeIfAbsent(preferredReader,k-newArrayList()).add(split);}// 下发分配for(varentry:readerAssignments.entrySet()){context.assignSplits(newSplitsAssignment(entry.getValue(),entry.getKey()));}}}数据读取ReaderpublicclassFlussSourceReaderimplementsSourceReaderRowData,FlussSourceSplit{OverridepublicvoidpollNext(ReaderOutputRowDataoutput){for(FlussSourceSplitsplit:assignedSplits){// 1. ★ 连接到 Split 对应的 TabletServerLogScannerscannergetOrCreateScanner(split);// 2. 读取一批 Arrow RecordBatchArrowRecordBatchbatchscanner.nextBatch();if(batch!null){// 3. ★ 列裁剪通过 projectedColumns 参数ArrowRecordBatchprojectedbatch.project(projectedColumns);// 4. Arrow → Flink RowData 转换for(inti0;iprojected.getRowCount();i){RowDatarowconvertToRowData(projected,i);output.collect(row);}// 5. 更新 Checkpoint offsetsplit.setCurrentOffset(scanner.getCurrentOffset());}}}}笔记 9.3FlussSink —— Sink 实现与 Exactly-Once文件FlussSink.java、FlussSinkWriter.java、FlussSinkCommitter.javaSink WriterpublicclassFlussSinkWriterimplementsSinkWriterRowData{privatefinalMapInteger,LogWriterbucketWriters;// BucketId → WriterOverridepublicvoidwrite(RowDatarow,Contextcontext){// 1. 确定分桶intbucketIdbucketingFunction.getBucket(row,numBuckets);// 2. 获取或创建对应 Bucket 的 WriterLogWriterwriterbucketWriters.computeIfAbsent(bucketId,id-createWriter(tablePath,partitionId,id));// 3. 序列化并写入ArrowRecordBatchbatchserializer.serialize(Collections.singletonList(row));writer.write(batch);}Overridepublicvoidflush(booleanendOfInput){// Flush 所有 pending 的 Batchfor(LogWriterwriter:bucketWriters.values()){writer.flush();}}}Two-Phase CommitExactly-Once 保证publicclassFlussSinkCommitterimplementsSinkCommitter{/** * ★ Phase 1: PrepareCheckpoint 触发时 * 将所有 Writer 的当前 offset 保存为 pending commit */publicListCommitRequestprepareCommit(){ListCommitRequestcommitsnewArrayList();for(varentry:bucketWriters.entrySet()){intbucketIdentry.getKey();LogWriterwriterentry.getValue();commits.add(newCommitRequest(tablePath,partitionId,bucketId,writer.getCurrentOffset()// ★ 记录当前已写入的 offset));}returncommits;}/** * ★ Phase 2: Commit所有并行 Writer 的 checkpoint 都完成后 * 将所有 pending commit 标记为已完成 */publicvoidcommit(ListCommitRequestcommits){for(CommitRequestreq:commits){// 通知 Fluss Server这批数据已成功写入并 Checkpoint// Server 端推进 Committed OffsetadminClient.commitOffset(req.tablePath,req.partitionId,req.bucketId,req.offset);}}/** * ★ 故障恢复从最近 Checkpoint 恢复 * Sink 自动从上次 Committed Offset 继续写入 * 不会产生重复数据因为 Checkpoint 前的数据已确认为 Committed */}笔记 9.4FlussLookupFunction —— Lookup Join文件FlussLookupFunction.javapublicclassFlussLookupFunctionextendsTableFunctionRowData{privatefinalFlussConnectionconnection;privatefinalCacheRowData,RowDatalookupCache;// ★ LRU 缓存/** * Flink 每来一条主表数据调用一次 eval * 对应 SQL: LEFT JOIN dim_table FOR SYSTEM_TIME AS OF o.time AS d ON o.key d.key */publicvoideval(Object...joinKeys){RowDatakeyGenericRowData.of(joinKeys);// 1. ★ 先查本地 LRU 缓存减少网络调用RowDatacachedlookupCache.getIfPresent(key);if(cached!null){collect(cached);return;}// 2. 缓存未命中 → 向 Fluss 发起 PK Lookupbyte[]lookupKeyserializeKey(joinKeys);byte[]resultconnection.pointLookup(FileSystemTablePath.of(dimTable),lookupKey);if(result!null){RowDatarowdeserializeRow(result);lookupCache.put(key,row);// 写入缓存collect(row);}// 维表中没有匹配的记录 → LEFT JOIN 只输出左表数据}/** * ★ 缓存配置 */publicstaticclassLookupCacheConfig{privatefinalintmaxRows;// 最大缓存行数默认 10000privatefinalDurationttl;// 缓存过期时间默认 10 分钟publicCacheRowData,RowDatacreateCache(){returnCaffeine.newBuilder().maximumSize(maxRows).expireAfterWrite(ttl).recordStats()// 记录缓存命中率.build();}}}阅读小结已理解尚未深入✅ FlussCatalog 如何将 Fluss Schema 转为 Flink Schema⬜$changelog和$binlog虚拟表的实现✅ SourceEnumerator 的 Split 发现和本地优先分配⬜ 动态分区发现运行时新增分区✅ Sink 的 Two-Phase Commit Exactly-Once 保证⬜SinkCommitter的 Globally Committed 生命周期✅ LookupFunction 的 LRU 缓存和点查询实现⬜ Lookup Join 对 Flink 执行计划的优化影响下一步功能点 10——DeltaJoinOperator、JoinStateStore 的状态外部化实现。
返回列表