云计算百科
云计算领域专业知识百科平台

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、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,调试标配

    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");
    }
    }

    两个最容易踩的配置点:

  • deliveryGuarantee 选 EXACTLY_ONCE 时,Checkpoint 必须开启,否则事务永远无法提交,作业会卡住或报错;
  • setTransactionalIdPrefix 必须全局唯一(不同作业不同前缀),否则两个作业抢同一批事务 ID,Kafka 直接抛 ProducerFencedException。

  • 四、实战三: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。

    两阶段提交的完整时序:

  • 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 写法最简单,适合内部系统对接:

    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 快速实现 + 下游幂等兜底」。


    九、六个真实踩坑

  • 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 并行度要与下游容量对齐。

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » Flink基础之Sink API详解及代码实战演练:数据出口的最后一公里
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!