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

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

 更新時間: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序列化器,并通過加入自己的邏輯主動往指定的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())){
        	// 動態(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中消費數(shù)據(jù)
DataStreamSource<String> kafkaDataStream =
        env.addSource(new FlinkKafkaConsumer<>(TOPIC, new SimpleStringSchema(), props));
// 3. 針對流做處理 把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é)

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

相關(guān)文章

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

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

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

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

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

    SpringBoot?spring.factories加載時機分析

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

    maven+springboot打成jar包的方法

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

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

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

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

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

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

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

    java程序員如何編寫更好的單元測試的7個技巧

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

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

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

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

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

最新評論

望谟县| 黑河市| 红河县| 江津市| 海南省| 苏尼特左旗| 浦城县| 锡林浩特市| 汝南县| 长海县| 丹江口市| 屏东县| 从江县| 无极县| 永福县| 皮山县| 上饶县| 长春市| 兴山县| 铜梁县| 岑溪市| 扬州市| 孟津县| 双鸭山市| 盘锦市| 海安县| 明溪县| 龙泉市| 财经| 天柱县| 苍山县| 无锡市| 遵义市| 阿鲁科尔沁旗| 高青县| 时尚| 巩义市| 冀州市| 阿合奇县| 洛南县| 孝义市|