Flink的Sink算子是流处理管道中的最终操作节点,负责将处理后的数据输出到外部系统。以下是其核心要点:

核心功能与定位

Sink算子通过addSinksinkTo方法实现数据输出,是Flink作业的终点站,支持将结果写入文件系统、数据库、消息队列等外部存储。其设计需保障状态一致性,通过检查点机制确保故障恢复时的数据正确性。

主要类型与实现

  1. 文件系统Sink

    • 早期使用writeAsText/writeAsCsv(已弃用),并行度影响文件输出形式(单文件或目录)。
    • 推荐使用StreamingFileSink,支持行编码(forRowFormat)和批量编码(forBulkFormat),自动分桶存储,适应分布式环境。
  2. 数据库Sink

    • 通过JDBCOutputFormat或自定义RichSinkFunction实现MySQL等关系型数据库写入,需在open生命周期建立连接。
  3. 消息队列Sink

    • Kafka集成需配置生产者地址、主题及序列化器(如SimpleStringSchema),依赖flink-connector-kafka模块。
  4. 自定义Sink
    实现SinkFunction接口并重写invoke方法,可灵活对接任意外部系统。例如,通过RichSinkFunction复用连接资源。

关键机制

  • 二阶段提交协议:保障端到端精确一次(Exactly-Once)语义,协调Flink与外部系统的数据一致性。
  • 分桶策略:文件Sink默认按时间分桶(如每小时新桶),支持自定义分区规则。

版本演进

  • Flink 1.12前使用addSink,之后推荐sinkTo API,架构更清晰。
  • FileSink替代StreamingFileSink,统一批流写入接口。

例子

  1. 文件系统Sink
/**
 * Flink专门提供了一个流式文件系统的连接器:FileSink,为批处理和流处理提供了一个统一的Sink,它可以将分区文件写入Flink支持的文件系统。
 * FileSink支持行编码(Row-encoded)和批量编码(Bulk-encoded)格式。这两种不同的方式都有各自的构造器,可以直接调用FileSink的静态方法。
 * 行编码:FileSink.forRowFormat(bathPath,rowEncoder)
 * 批量编码:FileSink.forBulkFormat(bathPath,bulkWriterFactory)
 */
public class SinkFile {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 每个目录中,都有并行度个数的文件在写入
        env.setParallelism(2);
        // 必须开启checkpoint,否则一直都是.inprogress
        env.enableCheckpointing(2000, CheckpointingMode.EXACTLY_ONCE);


        SingleOutputStreamOperator<String> map = env.fromElements(123L,
                234L, 345L
        ).map(String::valueOf);
        // 输出到文件系统
        FileSink<String> fileSink = FileSink.<String>forRowFormat(new Path("input/"),new SimpleStringEncoder<>("UTF-8"))
                .withOutputFileConfig(
                        OutputFileConfig.builder()
                                .withPartPrefix("atguigu-")
                                .withPartSuffix(".log")
                                .build()
                )
                // 按目录分桶:如下,每小时一个目录
                .withBucketAssigner(new DateTimeBucketAssigner<>("yyyy-MM-dd HH", ZoneId.systemDefault()))
                // 文件的滚动策略:1分钟或者一秒
                .withRollingPolicy(
                        DefaultRollingPolicy.builder()
                                .withRolloverInterval(Duration.ofMinutes(1))
                                .withMaxPartSize(new MemorySize(1024*1024))
                                .build())
                .build();

                map.sinkTo(fileSink);

                env.execute();
    }
}

  1. 数据库Sink
public class SinkMySQL {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        SingleOutputStreamOperator<WaterSensor> sensorDS = env.socketTextStream("hadoop101", 7777)
                .map(new WaterSensorMapFunction());

        /**
         * TODO 写入mysql
         * 1.只能用老的sink写法:addsink
         * 2.JDBCSink的四个参数
         *  第一个参数:执行的sql,
         *  第二个参数:预编译sql,对占位符填充值
         *  第三个参数:执行选项--->攒批、重试
         *  第四个参数:连接选线--->url、用户名、密码
         *
         */
        SinkFunction<WaterSensor> jdbcSink = JdbcSink.sink(
                "insert into ws values(?,?,?)",
                new JdbcStatementBuilder<WaterSensor>() {
                    @Override
                    public void accept(PreparedStatement preparedStatement, WaterSensor waterSensor) throws SQLException {
                        preparedStatement.setString(1, waterSensor.getId());
                        preparedStatement.setLong(2,waterSensor.getTs());
                        preparedStatement.setInt(3,waterSensor.getVc());
                    }
                },
                JdbcExecutionOptions.builder()
                        .withMaxRetries(3)// 重试次数
                        .withBatchSize(100)// 批次的大小:条数
                        .withBatchIntervalMs(3000)
                        .build(),
                new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
                        .withUrl("")
                        .withUsername("root")
                        .withUsername("root")
                        .withPassword("123456")
                        .withConnectionCheckTimeoutSeconds(60)
                        .build()
        );

        sensorDS.addSink(jdbcSink);

        env.execute();
    }
}

  1. 消息队列Sink
/**
 * 1.添加Kafka连接器依赖
 * 2.启动kafka集群
 * 3.编写输出到kafka的示例代码
 */
public class SinkKafka {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        // 如果是精准一次,必须开启checkpoint
        env.enableCheckpointing(2000, CheckpointingMode.EXACTLY_ONCE);

        SingleOutputStreamOperator<String> sensorDS = env.socketTextStream("hadoop101",7777);
        /**
         * Kafka sinkL
         * 注意:如果要使用精准一次写入kafka,需要满足一下条件。
         * 1.开启checkpoint
         * 2.设置事务前缀
         * 3.设置事务超时时间:checkpoint间隔<事务超时时间<max的15分钟
         *
         */
        KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
                .setBootstrapServers("")
                // 指定序列化器,指定topic名称,具体的序列化
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.<String>builder()
                                .setTopic("ws")
                                .setValueSerializationSchema(new SimpleStringSchema())
                                .build()
                )
                // 写到kafka的一致性级别:精准一次、至少一次
                .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
                // 如果是精准一次,必须设置事务的前缀
                .setTransactionalIdPrefix("atguigu-")
                // 如果是精准一次,必须设置事务超时时间:大于checkpoint间隔,小于max15分钟
                .setProperty(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG,10*60*1000+"")
                .build();

        sensorDS.sinkTo(kafkaSink);
        env.execute();

    }
}
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐