摘要
讲透 Flink Sink API 的完整体系:print 调试输出、KafkaSink / FileSink / JdbcSink 生产级 Connector、新 Sink API 的 SinkWriter 与 SinkCommitter 两阶段提交机制、输出一致性三种语义(at-least-once / 幂等 / exactly-once)的选型,附 5 个可直接运行的实战案例与 6 个真实踩坑点。
关键词
Flink、Sink API、KafkaSink、FileSink、JdbcSink、两阶段提交、exactly-once、幂等写、SinkWriter、SinkCommitter
Source 篇讲了数据怎么进来,Transformation 篇讲了数据怎么加工,这一篇收官:数据怎么出去。Sink 是作业的最后一公里,也是「exactly-once」这个承诺最终落地的地方——很多作业状态算得完全正确,却因为 Sink 配置错了,把数据写重、写丢、写慢。
这篇按「调试 → 生产 → 原理 → 语义」的顺序讲透 Sink:5 个能直接跑的案例,外加一个大多数教程不讲的关键点——两阶段提交到底在提交什么。
一、Sink API 分类全景
和 Source 对称,Sink 也有三类入口:
- 环境方法:
print()、printToErr()、writeAsText()(旧)——输出到标准输出或文件,只适合调试。 - Connector Sink:
KafkaSink、FileSink、JdbcSink、PulsarSink等官方实现——生产主力。 - 自定义 Sink:旧 API
SinkFunction/RichSinkFunction,新 API 的Sink + SinkWriter + SinkCommitter——兜底手段。
选型路径与 Source 一致:先找官方 Connector,没有才自定义。
同样要注意 API 演进:旧的env.addSink(new SinkFunction...)和FlinkKafkaProducer已废弃,新项目用stream.sinkTo(Sink)统一入口。新 Sink API 最大的升级是把「写数据」和「提交」拆成两个角色——这正是 exactly-once 的机制基础,后面细讲。
二、实战一:print,调试标配
importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;publicclassPrintSinkDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();DataStream<String>stream=env.fromElements("a","b","c");// 调试输出:每并行实例一行,格式 [subtaskIndex] valuestream.print("debug-tag");// 错误流输出(红色字体,方便区分)stream.printToErr();env.execute("print-sink-demo");}}print的注意事项:
- 输出格式
1> value,前面的数字是子任务编号,调试并行度问题很有用; - 它输出到TaskManager 的标准输出,集群模式下要在 TaskManager 日志里找,不是提交机;
- 生产环境禁止用 print 当正式 Sink——没有写入保障,数据会丢。
三、实战二:KafkaSink,生产最常用
importorg.apache.flink.api.common.serialization.SimpleStringSchema;importorg.apache.flink.connector.base.DeliveryGuarantee;importorg.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;importorg.apache.flink.connector.kafka.sink.KafkaSink;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;publicclassKafkaSinkDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();// 注意:EXACTLY_ONCE 语义依赖 Checkpoint,必须开启env.enableCheckpointing(30_000);DataStream<String>stream=env.fromElements("{\"orderId\":1,\"amount\":99.5}");KafkaSink<String>sink=KafkaSink.<String>builder().setBootstrapServers("localhost:9092").setRecordSerializer(KafkaRecordSerializationSchema.builder().setTopic("orders-out")// 整个 String 作为 value 写出去.setValueSerializationSchema(newSimpleStringSchema()).build())// 三种语义:EXACTLY_ONCE / AT_LEAST_ONCE / NONE.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)// EXACTLY_ONCE 下需要事务前缀(同集群唯一).setTransactionalIdPrefix("flink-order-sink").build();stream.sinkTo(sink);env.execute("kafka-sink-demo");}}两个最容易踩的配置点:
deliveryGuarantee选 EXACTLY_ONCE 时,Checkpoint 必须开启,否则事务永远无法提交,作业会卡住或报错;setTransactionalIdPrefix必须全局唯一(不同作业不同前缀),否则两个作业抢同一批事务 ID,Kafka 直接抛ProducerFencedException。
四、实战三:FileSink,落盘到 HDFS/本地
importorg.apache.flink.api.common.serialization.SimpleStringEncoder;importorg.apache.flink.connector.file.sink.FileSink;importorg.apache.flink.core.fs.Path;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy;importjava.time.Duration;publicclassFileSinkDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(30_000);DataStream<String>stream=env.fromElements("line1","line2");FileSink<String>sink=FileSink.forRowFormat(newPath("hdfs:///data/flink-output"),newSimpleStringEncoder<String>("UTF-8"))// 滚动策略:128MB 或 60s 或 30s 无新数据 → 关闭当前文件.withRollingPolicy(DefaultRollingPolicy.builder().withMaxPartSize(128*1024*1024).withRolloverInterval(Duration.ofSeconds(60)).withInactivityInterval(Duration.ofSeconds(30)).build()).build();stream.sinkTo(sink);env.execute("file-sink-demo");}}FileSink 两个关键机制:
- 目录按时间分桶:默认按处理时间分桶(
2026-08-24--12),数据落到对应桶目录; - 两阶段文件管理:写入中的文件是
.in-progress后缀,Checkpoint 成功后才 rename 为正式文件——所以在作业运行中看到的.in-progress文件不代表数据已落库,Checkpoint 完成才算数。这是 FileSink 语义正确的核心。
五、实战四:JdbcSink,写 MySQL
importorg.apache.flink.connector.jdbc.JdbcConnectionOptions;importorg.apache.flink.connector.jdbc.JdbcExecutionOptions;importorg.apache.flink.connector.jdbc.JdbcSink;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;publicclassJdbcSinkDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();// 输入:orderId,amount,categoryDataStream<String>stream=env.fromElements("1001,99.5,book","1002,12.0,music");// 批量 + upsert:同订单重复写时按主键覆盖 → 幂等,重复处理无害stream.sinkTo(JdbcSink.sink(// SQL:ON DUPLICATE KEY UPDATE 是幂等关键"INSERT INTO order_stats(order_id, amount, category) VALUES (?, ?, ?) "+"ON DUPLICATE KEY UPDATE amount = VALUES(amount)",(ps,line)->{String[]p=line.split(",");ps.setString(1,p[0]);ps.setDouble(2,Double.parseDouble(p[1]));ps.setString(3,p[2]);},// 执行参数:批量 1000 条提交一次,重试 3 次JdbcExecutionOptions.builder().withBatchSize(1000).withMaxRetries(3).build(),newJdbcConnectionOptions.JdbcConnectionOptionsBuilder().withUrl("jdbc:mysql://localhost:3306/ods?useSSL=false").withDriverName("com.mysql.cj.jdbc.Driver").withUsername("root").withPassword("password").build()));env.execute("jdbc-sink-demo");}}JdbcSink 的语义真相要认清:MySQL 不支持跨实例事务,JdbcSink 做不到真正的两阶段提交,它靠的是「批量写入 +ON DUPLICATE KEY UPDATE幂等」来逼近 exactly-once——重复写同一行,结果不变。这就是「幂等写」路线,是数据库 Sink 的主流做法。
性能关键:withBatchSize必须配。默认是逐条提交,每条一个事务,TPS 低到没法用;配 1000 批量后吞吐提升一个数量级。
六、新 Sink API 架构:两阶段提交到底在提交什么
理解新 Sink API,关键是两个角色:
- SinkWriter(运行在 TaskManager,每并行实例一个):接收记录,序列化后写入外部系统的暂存区——Kafka 的未提交事务、FileSink 的
.in-progress文件、JDBC 的批量缓冲。 - SinkCommitter(Checkpoint 成功后触发):拿到 Writer 产出的Committable(待提交句柄),执行真正的提交——Kafka 事务 commit、文件 rename。
两阶段提交的完整时序:
- Checkpoint barrier 到达 → SinkWriter 停止写入,生成 Committable 并随 checkpoint 持久化(预提交);
- Checkpoint 全局完成 → JobManager 通知 SinkCommitter;
- SinkCommitter 用 Committable 提交事务/rename 文件(确认提交);
- 故障恢复时,未提交的事务回滚,已提交的不重复提交——不多不少,这就是 exactly-once。
注意一个边界:exactly-once 是「Flink 两阶段提交 + 外部系统支持事务」的组合拳。Kafka 支持事务可以;文件系统支持原子 rename 可以;MySQL 这类不支持分布式事务的,只能退而求其次走幂等。
七、输出一致性:三种语义怎么选
| 语义 | 含义 | 实现 | 适用 |
|---|---|---|---|
| at-least-once | 可能重复 | 写完即算 | 日志、监控(容忍重复) |
| 幂等写 | 重复无害 | upsert / 覆盖写 | 数仓、MySQL 宽表 |
| exactly-once | 不重不漏 | 两阶段提交事务 | 金融、强一致链路 |
工程代价从低到高,选型建议:
- 下游是 Kafka →
EXACTLY_ONCE(开启 checkpoint,吞吐损失可控); - 下游是 MySQL/HBase →幂等写(设计业务幂等键 + upsert),性价比最高;
- 下游只是日志/告警 → at-least-once 就够,别为用不上的一致性牺牲吞吐。
判断自己的场景:先问「重复写一条数据,下游会不会出事?」不会 → 幂等/at-least-once 都行;会 → 必须 exactly-once 或强幂等。
八、实战五:自定义 Sink(兜底)
官方没有现成 Connector 时才自定义。旧 API 的RichSinkFunction写法最简单,适合内部系统对接:
importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.functions.sink.RichSinkFunction;publicclassCustomSinkDemo{publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();env.fromElements("event-1","event-2").addSink(newRichSinkFunction<String>(){@Overridepublicvoidopen(org.apache.flink.configuration.Configurationparameters){// 生命周期方法:初始化连接(每并行实例调用一次)System.out.println("open connection");}@Overridepublicvoidinvoke(Stringvalue,Contextcontext){// 每来一条数据调用一次:写外部系统System.out.println("write: "+value);}@Overridepublicvoidclose(){// 作业结束时释放连接System.out.println("close connection");}});env.execute("custom-sink-demo");}}必须知道的坑:RichSinkFunction没有两阶段提交能力——故障重放时invoke会被重复调用,数据重复写。要么自己实现幂等(下游按业务键去重),要么升级到新 Sink API 的Sink + SinkWriter + SinkCommitter(代码量大幅上升,但语义完整)。生产对接自研系统时,多数团队选择「旧 API 快速实现 + 下游幂等兜底」。
九、六个真实踩坑
- print 当生产 Sink 用。print 输出到 TaskManager stdout,日志一滚就没了;生产必须接正式 Sink。本地调试完记得删掉。
- KafkaSink 的 EXACTLY_ONCE 没开 checkpoint。事务没有触发点,数据永远停在未提交状态;而且
transactionalIdPrefix不唯一会互相打架,报ProducerFencedException。 - FileSink 的
.in-progress文件误以为已落库。文件要等 checkpoint 成功后才 rename;下游消费这些文件时,要么读正式文件,要么配合 checkpoint 时机。反过来,不配 checkpoint 时 FileSink 文件永远不会变正式。 - JdbcSink 不配 batchSize。默认逐条提交,每条一个事务,几千 TPS 就是上限。配
withBatchSize(1000)后量级提升。 - 自定义 Sink 无幂等还开 exactly-once。
RichSinkFunction没有提交语义,上游 checkpoint 恢复后 invoke 必然重放;下游没有去重键,数据就重复了。要么下游幂等,要么别对外宣称精确一次。 - Sink 并行度不匹配下游容量。KafkaSink 并行度 > 目标 topic 分区数,多出来的实例空闲;JDBC Sink 并行度高但连接数也翻倍,可能打爆数据库连接池——Sink 并行度要与下游容量对齐。
Sink 是作业的最后一公里,也是语义承诺的兑现处。把「print 只调试」「官方 Connector 优先」「exactly-once = 两阶段提交 + 外部系统支持」「数据库用幂等写逼近精确一次」这四句话记住,配合「先问重复写会不会出事」的选型方法,Sink 环节就不会再掉链子。至此,Source → Transform → Sink 三件套全部讲完,一个完整的 Flink DataStream 作业从入口到出口的每一环都有了清晰认知。
闲;JDBC Sink 并行度高但连接数也翻倍,可能打爆数据库连接池——Sink 并行度要与下游容量对齐。