发布时间:2026/8/25 6:09:36
功能点 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 的状态外部化实现。

相关新闻

2026/8/25 6:09:36

大模型如何重构软件价值:从功能增强到架构重生的实战路径

1. 风暴来临:当“颠覆”不再是口号,而是生存抉择最近和几个做企业软件的朋友聊天,话题绕不开一个词:焦虑。这种焦虑不是来自竞争对手,也不是来自市场波动,而是来自一个更根本、更不可逆的趋势——大模型。过…

2026/8/25 6:04:36

音频内容智能审核:4倍速并行架构与工程实践全解析

1. 项目概述:当“耳朵经济”遇上内容安全最近几年,长音频内容(比如播客、有声书、在线课程、直播回放)的爆发式增长,让“耳朵经济”成了一个实实在在的风口。但随之而来的,是平台运营者一个头两个大的难题&…

2026/8/25 6:04:36

AI编程CLI配置管理工具:告别配置地狱,统一管理多模型环境

1. 项目概述:为什么我们需要一个AI编程CLI配置管理工具?如果你和我一样,每天的工作流里充斥着各种命令行工具,尤其是那些与AI编程相关的——比如调用不同模型的API、切换项目环境、管理一堆API密钥和配置文件——那你肯定对“配置…

2026/8/25 8:29:57

Linear Attention技术解析与面试准备指南

1. Linear Attention技术解析与面试准备指南在自然语言处理和计算机视觉领域,Attention机制已经成为现代深度学习模型的核心组件。然而传统Attention计算中的平方复杂度问题,一直是制约模型效率的瓶颈。Linear Attention作为突破这一限制的创新方法&…

2026/8/25 8:29:57

AI工程师转型实战:3个月高效学习路线与求职策略

1. 项目概述:AI转型的关键窗口期去年夏天,我辞去了做了五年的传统行业数据分析工作,决定All in AI领域。当时身边所有人都说我疯了——29岁转行、零AI项目经验、Python只会写基础脚本。但三个月后,我手握三家头部科技公司的AI工程…

2026/8/25 8:29:57

01背包问题详解:从采药题看动态规划本质

1. 这道题不是“采药”,是01背包问题的成人礼你第一次在信息学奥赛一本通里翻到“1290:采药”这页时,大概率正坐在教室后排,手边摊着一本翻得卷了边的《信息学奥赛一本通》,旁边堆着几本洛谷打印题单。老师刚讲完“递归…

2026/8/25 8:29:57

2026前端面试趋势与React 18、Vue 3核心考点解析

1. 2026前端面试趋势与核心能力解析2026年的前端技术栈正在经历新一轮洗牌。从各大厂最新招聘需求和社区技术风向来看,React 18和Vue 3的组合式API已成为标配技能,但仅掌握框架API早已不够。我在最近三个月参与了7场不同规模公司的技术面试,发…

2026/8/25 1:04:19

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/24 1:12:32

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/24 8:17:29

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/25 0:04:14

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory Meta Description:GetQzonehistory 是一个QQ空间历史说…

2026/8/25 0:04:14

洛谷 P7912:[CSP-J 2021 T4] 小熊的果篮 ← 双向链表

【题目来源】 https://www.luogu.com.cn/problem/P7912 【题目描述】 小熊的水果店里摆放着一排 n 个水果。每个水果只可能是苹果或桔子,从左到右依次用正整数 1,2,…,n 编号。连续排在一起的同一种水果称为一个“块”。小熊要把这一排水果挑到若干个果篮里&#x…

2026/8/24 13:42:17

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/24 18:13:48

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/25 1:08:14

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…