【问题标题】:Flink Push Row to kafkaFlink Push Row 到 kafka
【发布时间】:2023-03-11 09:10:02
【问题描述】:

我有一个 flink Row 和列名,这个Row 可以通过字段名或索引访问。我想使用 vanilla flink kafka producer 将它放入 JSON 中的 kafka 中。我该怎么做?目标 Json Schema 是否需要将其下沉到 kafka?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    您需要为您为流指定的 kafka 生产者提供架构。幸运的是,Flink 确实为您提供了可以根据需要进行修改的模式。如果您只想将您的对象作为 JSON 字符串发送,您可以执行以下操作:

    将您的对象流转换为字符串 JSON 流,如下所示:

    SingleOutputStreamOperator<String> jsons = dataStream.map(new MapFunction<Object, String>() {
        @Override
        public String map(Object value) throws Exception {
            // Gson creation can be put in a static utility method if you want to avoid recreating 
            Gson gson = new GsonBuilder().create();
            return gson.toJson(value);
        }
    });
    

    然后您可以按如下方式定义 Kafka Producer:

        public static FlinkKafkaProducer<String> getKafkaProducer(String topic) {
            String kafkaBootstrapServers = "localhost:9092";
            String kafkaGroup = "kafkaGroup";
    
            Properties propertiesProducer = new Properties();
            propertiesProducer.setProperty("bootstrap.servers", kafkaBootstrapServers);
            propertiesProducer.setProperty("group.id", kafkaGroup);
    
            SimpleStringSchema simpleStringSchema = new SimpleStringSchema() {
    
                public String deserialize(byte[] message) {
                    return message == null ? null : super.deserialize(message);
                }
            };
            return new FlinkKafkaProducer(topic, simpleStringSchema, propertiesProducer);
        }
    

    最后为流指定 kafka 生产者并使用它来接收您的 JSON 字符串消息:

    jsons
        .addSink(FlinkUtils.getKafkaProducer("outputTopic))
        .name("JSON messages sink");
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-08-13
      • 2019-08-01
      • 2023-03-07
      • 2018-10-17
      • 2021-05-15
      • 2017-01-23
      • 2017-06-27
      相关资源
      最近更新 更多