flink api-datastream api-sink算子
·
Flink的Sink算子是流处理管道中的最终操作节点,负责将处理后的数据输出到外部系统。以下是其核心要点:
核心功能与定位
Sink算子通过addSink或sinkTo方法实现数据输出,是Flink作业的终点站,支持将结果写入文件系统、数据库、消息队列等外部存储。其设计需保障状态一致性,通过检查点机制确保故障恢复时的数据正确性。
主要类型与实现
-
文件系统Sink
- 早期使用
writeAsText/writeAsCsv(已弃用),并行度影响文件输出形式(单文件或目录)。 - 推荐使用
StreamingFileSink,支持行编码(forRowFormat)和批量编码(forBulkFormat),自动分桶存储,适应分布式环境。
- 早期使用
-
数据库Sink
- 通过JDBCOutputFormat或自定义
RichSinkFunction实现MySQL等关系型数据库写入,需在open生命周期建立连接。
- 通过JDBCOutputFormat或自定义
-
消息队列Sink
- Kafka集成需配置生产者地址、主题及序列化器(如
SimpleStringSchema),依赖flink-connector-kafka模块。
- Kafka集成需配置生产者地址、主题及序列化器(如
-
自定义Sink
实现SinkFunction接口并重写invoke方法,可灵活对接任意外部系统。例如,通过RichSinkFunction复用连接资源。
关键机制
- 二阶段提交协议:保障端到端精确一次(Exactly-Once)语义,协调Flink与外部系统的数据一致性。
- 分桶策略:文件Sink默认按时间分桶(如每小时新桶),支持自定义分区规则。
版本演进
- Flink 1.12前使用
addSink,之后推荐sinkToAPI,架构更清晰。 FileSink替代StreamingFileSink,统一批流写入接口。
例子
- 文件系统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();
}
}
- 数据库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();
}
}
- 消息队列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();
}
}
更多推荐



所有评论(0)