最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Flink實(shí)現(xiàn)往Kafka中多個(gè)topic發(fā)送消息

 更新時(shí)間:2025年10月13日 09:44:47   作者:暴走的Aluuubbarrrr  
文章介紹了使用Flink 1.13.2 和 Kafka 2.6.2 從Kafka讀取數(shù)據(jù),并根據(jù)邏輯將數(shù)據(jù)分配到不同的topic,提到了需要重寫FlinkKafka的Key序列化器,并加入自定義邏輯以發(fā)送消息到指定的topic,文中還包括了配置Kafka信息和如何連接FlinkKafka的步驟
  • Flink 1.13.2
  • Kafka 2.6.2

思路與環(huán)境

從kafka中讀取數(shù)據(jù) 根據(jù)邏輯判斷分配到不同的topic中去

需要重寫Flink Kafka的Key序列化器,并通過(guò)加入自己的邏輯主動(dòng)往指定的topic發(fā)送消息。

首先配置Kafka信息

Properties props = new Properties();
props.put("bootstrap.servers","10.116.0.16:9092");
props.put("acks", "all");
props.put("retries", 1);
props.put("batch.size", 16384);
props.put("linger.ms", 1);
props.put("buffer.memory", 33554432);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");

注意這里key序列化器和value序列化器都為StringSerializer

Flink Kafka連接器

FlinkKafkaProducer<FlinkJobBO> fkProducer =
                new FlinkKafkaProducer<>("", new MyKeySerialization(), props, FlinkKafkaProducer.Semantic.EXACTLY_ONCE);
        

其中MyKeySerialization便是重寫的key序列化器

自定義序列化器

public class MyKeySerialization implements KafkaSerializationSchema<FlinkJobBO> {
    String topic;
    public MyKeySerialization(String topic){
        this.topic = topic;
    }
    public MyKeySerialization(){
    }

	// 注意:都是byte[]類型,所以我們要重新指定新的序列化器
    @Override
    public ProducerRecord<byte[], byte[]> serialize(FlinkJobBO flinkJobBO, @Nullable Long aLong) {
    	// 根據(jù)自身的邏輯條件
        JsonUtils.setObjectMapper(new ObjectMapper());
        if("1".equals(flinkJobBO.getApiModel())){
        	// 動(dòng)態(tài)生成topic
            return new ProducerRecord<>("topic-"+flinkJobBO.getGroupId(), JsonUtils.toJson(flinkJobBO).getBytes(StandardCharsets.UTF_8));
        }
        return new ProducerRecord<>("", "".getBytes(StandardCharsets.UTF_8));
    }
}

需要把

props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");

替換為

props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");

完整代碼

// 創(chuàng)建Flink Stream執(zhí)行環(huán)境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// 1、設(shè)置默認(rèn)topic
String TOPIC = "TEST";
// 2. 從kafka獲取流數(shù)據(jù)
Properties props = new Properties();
props.put("bootstrap.servers","10.116.0.16:9092");
props.put("acks", "all");
props.put("retries", 1);
props.put("batch.size", 16384);
props.put("linger.ms", 1);
props.put("buffer.memory", 33554432);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
// 從kafka中消費(fèi)數(shù)據(jù)
DataStreamSource<String> kafkaDataStream =
        env.addSource(new FlinkKafkaConsumer<>(TOPIC, new SimpleStringSchema(), props));
// 3. 針對(duì)流做處理 把string轉(zhuǎn)成bo 主流
DataStream<FlinkJobBO> ds = kafkaDataStream
        .map((MapFunction<String, FlinkJobBO>) s -> JsonUtils.toBean(s, FlinkJobBO.class));
// 修改value序列化器
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
// 4.1.1 自定義序列化器 分配topic *****
FlinkKafkaProducer<FlinkJobBO> fkProducer =
        new FlinkKafkaProducer<>("", new MyKeySerialization(), props, FlinkKafkaProducer.Semantic.EXACTLY_ONCE);
fkProducer.setLogFailuresOnly(false);
ds.addSink(fkProducer);
env.execute();

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • 10k+點(diǎn)贊的 SpringBoot 后臺(tái)管理系統(tǒng)教程詳解

    10k+點(diǎn)贊的 SpringBoot 后臺(tái)管理系統(tǒng)教程詳解

    這篇文章主要介紹了10k+點(diǎn)贊的 SpringBoot 后臺(tái)管理系統(tǒng)教程詳解,本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-01-01
  • 淺析Java語(yǔ)言中狀態(tài)模式的優(yōu)點(diǎn)

    淺析Java語(yǔ)言中狀態(tài)模式的優(yōu)點(diǎn)

    狀態(tài)模式允許對(duì)象在內(nèi)部狀態(tài)改變時(shí)改變它的行為,對(duì)象看起來(lái)好像修改了它的類。這個(gè)模式將狀態(tài)封裝成獨(dú)立的類,并將動(dòng)作委托到 代表當(dāng)前狀態(tài)的對(duì)象,我們知道行為會(huì)隨著內(nèi)部狀態(tài)而改變
    2023-02-02
  • SpringBoot?spring.factories加載時(shí)機(jī)分析

    SpringBoot?spring.factories加載時(shí)機(jī)分析

    這篇文章主要為大家介紹了SpringBoot?spring.factories加載時(shí)機(jī)分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-03-03
  • maven+springboot打成jar包的方法

    maven+springboot打成jar包的方法

    這篇文章主要介紹了maven+springboot打成jar包的方法,非常不錯(cuò),具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2018-10-10
  • 詳解關(guān)于eclipse中使用jdk15對(duì)應(yīng)javafx15的配置問(wèn)題總結(jié)

    詳解關(guān)于eclipse中使用jdk15對(duì)應(yīng)javafx15的配置問(wèn)題總結(jié)

    這篇文章主要介紹了詳解關(guān)于eclipse中使用jdk15對(duì)應(yīng)javafx15的配置問(wèn)題總結(jié),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2020-11-11
  • 使用Sentinel實(shí)現(xiàn)流控和服務(wù)降級(jí)的代碼示例

    使用Sentinel實(shí)現(xiàn)流控和服務(wù)降級(jí)的代碼示例

    Sentinel是面向分布式、多語(yǔ)言異構(gòu)化服務(wù)架構(gòu)的流量治理組件,本文將詳細(xì)為大家介紹如何使用Sentinel實(shí)現(xiàn)流控和服務(wù)降級(jí),文中有相關(guān)的代碼示例,需要的朋友可以參考下
    2023-05-05
  • Java 泛型通配符 <? extends> vs <? super> 實(shí)戰(zhàn)場(chǎng)景

    Java 泛型通配符 <? extends> vs <? s

    本文解析Java泛型通配符<? extends T>和<? super T>的核心區(qū)別與應(yīng)用場(chǎng)景,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2025-12-12
  • java程序員如何編寫更好的單元測(cè)試的7個(gè)技巧

    java程序員如何編寫更好的單元測(cè)試的7個(gè)技巧

    測(cè)試是開發(fā)的一個(gè)非常重要的方面,可以在很大程度上決定一個(gè)應(yīng)用程序的命運(yùn)。良好的測(cè)試可以在早期捕獲導(dǎo)致應(yīng)用程序崩潰的問(wèn)題,但較差的測(cè)試往往總是導(dǎo)致故障和停機(jī)。本文主要介紹java程序員編寫更好的單元測(cè)試的7個(gè)技巧。下面跟著小編一起來(lái)看下吧
    2017-03-03
  • 詳解Spring注入集合(數(shù)組、List、Map、Set)類型屬性

    詳解Spring注入集合(數(shù)組、List、Map、Set)類型屬性

    這篇文章主要介紹了詳解Spring注入集合(數(shù)組、List、Map、Set)類型屬性,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2021-01-01
  • Java多線程實(shí)現(xiàn)阻塞隊(duì)列的示例代碼

    Java多線程實(shí)現(xiàn)阻塞隊(duì)列的示例代碼

    本文主要介紹了Java多線程實(shí)現(xiàn)阻塞隊(duì)列的示例代碼,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2024-12-12

最新評(píng)論

泾源县| 梁平县| 大渡口区| 若尔盖县| 繁峙县| 乐陵市| 东至县| 榆社县| 盘锦市| 宁陕县| 利津县| 安龙县| 公安县| 体育| 栾城县| 铜梁县| 福泉市| 卢龙县| 张家川| 青冈县| 哈密市| 酒泉市| 昌都县| 望都县| 贞丰县| 五台县| 高尔夫| 翼城县| 万载县| 阳西县| 高陵县| 德格县| 高州市| 昌乐县| 金寨县| 秭归县| 天津市| 会同县| 七台河市| 徐闻县| 彭山县|