开发者社区 > 大数据与机器学习 > 实时计算 Flink > 正文

1.15版本的flink往kafka多个topic写消息,是自定义实现KafkaSerializat

1.15版本的flink往kafka多个topic写消息,是自定义实现KafkaSerializationSchema吗?

展开
收起
wenti 2023-02-14 13:58:11 379 0
1 条回答
写回答
取消 提交回答
  • 是的,在Flink 1.15版本中向多个Kafka topic写入消息需要实现自定义的KafkaSerializationSchema。用于将消息转换为Kafka可以识别的二进制数据,以写入Kafka topic。实现的方法如下: public class MyKafkaSerializationSchema implements KafkaSerializationSchema<Tuple2<String, Integer>> {

    @Override
    public ProducerRecord<byte[], byte[]> serialize(Tuple2<String, Integer> element, @Nullable Long timestamp) {
        return new ProducerRecord<>(element.f0, element.f1.toString().getBytes());
    }
    

    }

    然后在使用FlinkKafkaProducer时,传入自定义的MyKafkaSerializationSchema即可——该回答整理自钉群“【③群】Apache Flink China社区”

    2023-02-14 16:38:41
    赞同 展开评论 打赏

实时计算Flink版是阿里云提供的全托管Serverless Flink云服务,基于 Apache Flink 构建的企业级、高性能实时大数据处理系统。提供全托管版 Flink 集群和引擎,提高作业开发运维效率。

相关产品

  • 实时计算 Flink版
  • 热门讨论

    热门文章

    相关电子书

    更多
    Java Spring Boot开发实战系列课程【第16讲】:Spring Boot 2.0 实战Apache Kafka百万级高并发消息中间件与原理解析 立即下载
    MaxCompute技术公开课第四季 之 如何将Kafka数据同步至MaxCompute 立即下载
    消息队列kafka介绍 立即下载