Flink如何从SingleOutputStreamOperator<ObjectA>流中提取ObjectB写入Kafka
Flink提取ObjectA流中嵌套ObjectB下沉Kafka实现方案
你的需求完全可以通过FlatMap实现,这也是Flink处理一对多数据流转换的标准方案,具体实现如下:
1. 实现打平逻辑
自定义FlatMapFunction,将每个ObjectA实例中的List
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.util.Collector; // 泛型分别对应输入类型ObjectA、输出类型ObjectB public class ObjectBExtractor implements FlatMapFunction<ObjectA, ObjectB> { @Override public void flatMap(ObjectA objectA, Collector<ObjectB> collector) throws Exception { // 增加判空逻辑,避免空指针导致作业崩溃 if (objectA.objBList != null && !objectA.objBList.isEmpty()) { for (ObjectB objB : objectA.objBList) { collector.collect(objB); } } } }
调用FlatMap得到ObjectB类型的流:
SingleOutputStreamOperator<ObjectB> objBStream = objAStream.flatMap(new ObjectBExtractor());
2. 下沉ObjectB流到Kafka
这里以Flink 1.14+版本的官方Kafka连接器为例,示例代码如下:
import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.kafka.clients.producer.ProducerConfig; import com.fasterxml.jackson.databind.ObjectMapper; // 构造Kafka Sink KafkaSink<String> kafkaSink = KafkaSink.<String>builder() .setBootstrapServers("your_kafka_broker_address:9092") .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic("your_target_topic_name") .setValueSerializationSchema(new SimpleStringSchema()) .build() ) .setConfig(ProducerConfig.ACKS_CONFIG, "all") .build(); // 将ObjectB序列化为JSON字符串后写入Kafka objBStream.map(objB -> { ObjectMapper objectMapper = new ObjectMapper(); return objectMapper.writeValueAsString(objB); }).sinkTo(kafkaSink);
可选简化写法
如果使用Flink 1.13及以上版本,也可以用map+flatten的组合写法,本质和FlatMap逻辑一致,代码更简洁:
import java.util.Collection; SingleOutputStreamOperator<ObjectB> objBStream = objAStream .map(objectA -> objectA.objBList) .flatMap(Collection::stream);
注意事项
- 如果使用旧版本Flink的
FlinkKafkaProducerAPI,替换为对应版本的Sink实现即可 - ObjectB的序列化方式可以根据业务需求调整,比如使用Avro、Protobuf等结构化序列化方案
- 生产环境建议将ObjectMapper实例声明为静态变量,避免重复创建对象带来的性能开销
内容的提问来源于stack exchange,提问作者chandana.ranasinghe
相关产品推荐
相关产品推荐

