【发布时间】:2020-09-09 03:41:56
【问题描述】:
现在我的 flink 代码正在处理一个文件,并在 1 个分区的 kafka 主题上下沉数据。
现在我有一个有 2 个分区的主题,我希望 flink 代码使用 DefaultPartitioner 在这 2 个分区上接收数据。
你能帮我解决这个问题吗?
这是我当前代码的sn-p代码:
DataStream<String> speStream = inputStream..map(new MapFunction<Row, String>(){....}
Properties props = Producer.getProducerConfig(propertiesFilePath);
speStream.addSink(new FlinkKafkaProducer011(kafkaTopicName, new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema()), props, FlinkKafkaProducer011.Semantic.EXACTLY_ONCE));
【问题讨论】:
标签: apache-flink