【问题标题】:Send in multiple topic kafka sink with flink使用 flink 发送多个主题 kafka sink
【发布时间】:2018-11-04 20:33:23
【问题描述】:

我有一个像这样的数据流:

DataStream[myTuple(topic, value)]

我想在相关主题中发送一个特定值。

所以我尝试做这样的事情:

new FlinkKafkaProducer010[myTuple](
  "default_topic",
  new KeyedSerializationSchema[myTuple](){
    override def getTargetTopic(element: myTuple): String = element.topic
    override def serializeKey(element: myTuple): Array[Byte] = null
    override def serializeValue(element: myTuple): Array[Byte] = new SimpleStringSchema().serialize(element.value)
  },
  properties)

但它不起作用,我有这个警告:

WARN  org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducerBase  - Overwriting the 'key.serializer' is not recommended
WARN  org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducerBase  - Overwriting the 'value.serializer' is not recommended

我不知道该怎么做,换一种方式。 谢谢你的帮助。

【问题讨论】:

    标签: scala apache-kafka apache-flink kafka-producer-api flink-streaming


    【解决方案1】:

    您可能在属性中设置了key.serializervalue.serializer。你不应该这样做,因为这样你会覆盖 Flink 内部使用的序列化程序(ByteArraySerializers)。删除这些属性,您的代码应该可以工作。

    【讨论】:

    • 在属性中我只添加bootstrap.servers
    • 请彻底检查,因为警告信息表明您确实...
    • 我解决了我的问题:问题是我尝试为每个条目打开一个 kafka 连接
    猜你喜欢
    • 1970-01-01
    • 2022-06-30
    • 1970-01-01
    • 2020-12-17
    • 2018-12-31
    • 2020-03-10
    • 1970-01-01
    • 2020-09-18
    • 2017-04-07
    相关资源
    最近更新 更多