摘要
讲透 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 也有三类入口:

选型路径与 Source 一致:先找官方 Connector,没有才自定义。
同样要注意 API 演进:旧的 env.addSink(new SinkFunction…) 和 FlinkKafkaProducer 已废弃,新项目用 stream.sinkTo(Sink) 统一入口。新 Sink API 最大的升级是把「写数据」和「提交」拆成两个角色——这正是 exactly-once 的机制基础,后面细讲。
二、实战一:print,调试标配
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class PrintSinkDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> stream = env.fromElements("a", "b", "c");
// 调试输出:每并行实例一行,格式 [subtaskIndex] value
stream.print("debug-tag");
// 错误流输出(红色字体,方便区分)
stream.printToErr();
env.execute("print-sink-demo");
}
}
print 的注意事项:
- 输出格式 1> value,前面的数字是子任务编号,调试并行度问题很有用;
- 它输出到 TaskManager 的标准输出,集群模式下要在 TaskManager 日志里找,不是提交机;
- 生产环境禁止用 print 当正式 Sink——没有写入保障,数据会丢。
三、实战二:KafkaSink,生产最常用
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.base.DeliveryGuarantee;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class KafkaSinkDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = 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(new SimpleStringSchema())
.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");
}
}
两个最容易踩的配置点:
四、实战三:FileSink,落盘到 HDFS/本地
import org.apache.flink.api.common.serialization.SimpleStringEncoder;
import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.core.fs.Path;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy;
import java.time.Duration;
public class FileSinkDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(30_000);
DataStream<String> stream = env.fromElements("line1", "line2");
FileSink<String> sink = FileSink.forRowFormat(
new Path("hdfs:///data/flink-output"),
new SimpleStringEncoder<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
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class JdbcSinkDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 输入:orderId,amount,category
DataStream<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(),
new JdbcConnectionOptions.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。
两阶段提交的完整时序:
注意一个边界:exactly-once 是「Flink 两阶段提交 + 外部系统支持事务」的组合拳。Kafka 支持事务可以;文件系统支持原子 rename 可以;MySQL 这类不支持分布式事务的,只能退而求其次走幂等。
七、输出一致性:三种语义怎么选

| at-least-once | 可能重复 | 写完即算 | 日志、监控(容忍重复) |
| 幂等写 | 重复无害 | upsert / 覆盖写 | 数仓、MySQL 宽表 |
| exactly-once | 不重不漏 | 两阶段提交事务 | 金融、强一致链路 |
工程代价从低到高,选型建议:
判断自己的场景:先问「重复写一条数据,下游会不会出事?」不会 → 幂等/at-least-once 都行;会 → 必须 exactly-once 或强幂等。
八、实战五:自定义 Sink(兜底)
官方没有现成 Connector 时才自定义。旧 API 的 RichSinkFunction 写法最简单,适合内部系统对接:
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
public class CustomSinkDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.fromElements("event-1", "event-2")
.addSink(new RichSinkFunction<String>() {
@Override
public void open(org.apache.flink.configuration.Configuration parameters) {
// 生命周期方法:初始化连接(每并行实例调用一次)
System.out.println("open connection");
}
@Override
public void invoke(String value, Context context) {
// 每来一条数据调用一次:写外部系统
System.out.println("write: " + value);
}
@Override
public void close() {
// 作业结束时释放连接
System.out.println("close connection");
}
});
env.execute("custom-sink-demo");
}
}
必须知道的坑:RichSinkFunction 没有两阶段提交能力——故障重放时 invoke 会被重复调用,数据重复写。要么自己实现幂等(下游按业务键去重),要么升级到新 Sink API 的 Sink + SinkWriter + SinkCommitter(代码量大幅上升,但语义完整)。生产对接自研系统时,多数团队选择「旧 API 快速实现 + 下游幂等兜底」。
九、六个真实踩坑
Sink 是作业的最后一公里,也是语义承诺的兑现处。把「print 只调试」「官方 Connector 优先」「exactly-once = 两阶段提交 + 外部系统支持」「数据库用幂等写逼近精确一次」这四句话记住,配合「先问重复写会不会出事」的选型方法,Sink 环节就不会再掉链子。至此,Source → Transform → Sink 三件套全部讲完,一个完整的 Flink DataStream 作业从入口到出口的每一环都有了清晰认知。
闲;JDBC Sink 并行度高但连接数也翻倍,可能打爆数据库连接池——Sink 并行度要与下游容量对齐。
网硕互联帮助中心







评论前必须登录!
注册