ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink基础之Sink API详解及代码实战演练:数据出口的最后一公里

Flink基础之Sink API详解及代码实战演练:数据出口的最后一公里 摘要讲透 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、SinkCommitterSource 篇讲了数据怎么进来Transformation 篇讲了数据怎么加工这一篇收官数据怎么出去。Sink 是作业的最后一公里也是「exactly-once」这个承诺最终落地的地方——很多作业状态算得完全正确却因为 Sink 配置错了把数据写重、写丢、写慢。这篇按「调试 → 生产 → 原理 → 语义」的顺序讲透 Sink5 个能直接跑的案例外加一个大多数教程不讲的关键点——两阶段提交到底在提交什么。一、Sink API 分类全景和 Source 对称Sink 也有三类入口环境方法print()、printToErr()、writeAsText()旧——输出到标准输出或文件只适合调试。Connector SinkKafkaSink、FileSink、JdbcSink、PulsarSink等官方实现——生产主力。自定义 Sink旧 APISinkFunction/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{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();DataStreamStringstreamenv.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{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// 注意EXACTLY_ONCE 语义依赖 Checkpoint必须开启env.enableCheckpointing(30_000);DataStreamStringstreamenv.fromElements({\orderId\:1,\amount\:99.5});KafkaSinkStringsinkKafkaSink.Stringbuilder().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必须全局唯一不同作业不同前缀否则两个作业抢同一批事务 IDKafka 直接抛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{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(30_000);DataStreamStringstreamenv.fromElements(line1,line2);FileSinkStringsinkFileSink.forRowFormat(newPath(hdfs:///data/flink-output),newSimpleStringEncoderString(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写 MySQLimportorg.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{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();// 输入orderId,amount,categoryDataStreamStringstreamenv.fromElements(1001,99.5,book,1002,12.0,music);// 批量 upsert同订单重复写时按主键覆盖 → 幂等重复处理无害stream.sinkTo(JdbcSink.sink(// SQLON DUPLICATE KEY UPDATE 是幂等关键INSERT INTO order_stats(order_id, amount, category) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE amount VALUES(amount),(ps,line)-{String[]pline.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?useSSLfalse).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 的批量缓冲。SinkCommitterCheckpoint 成功后触发拿到 Writer 产出的Committable待提交句柄执行真正的提交——Kafka 事务 commit、文件 rename。两阶段提交的完整时序Checkpoint barrier 到达 → SinkWriter 停止写入生成 Committable 并随 checkpoint 持久化预提交Checkpoint 全局完成 → JobManager 通知 SinkCommitterSinkCommitter 用 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{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.fromElements(event-1,event-2).addSink(newRichSinkFunctionString(){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 并行度要与下游容量对齐。
返回列表