日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

指定开始_Flink-Kafka指定offset的五种方式

發布時間:2023/12/10 编程问答 31 豆豆
生活随笔 收集整理的這篇文章主要介紹了 指定开始_Flink-Kafka指定offset的五种方式 小編覺得挺不錯的,現在分享給大家,幫大家做個參考.

默認:從topic中指定的group上次消費的位置開始消費。

所以必須配置group.id參數從消費者組提交的偏移量開始讀取分區(kafka或zookeeper中)。如果找不到分區的偏移量,auto.offset.reset將使用屬性中的設置。如果是默認行為(setStartFromGroupOffsets),那么任務從檢查點重啟,按照重啟前的offset進行消費,如果直接重啟不從檢查點重啟并且group.id不變,程序會按照上次提交的offset的位置繼續消費。如果group.id改變了,則程序按照auto.offset.reset設置的屬性進行消費。但是如果程序帶有狀態的算子,還是建議使用檢查點重啟。

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(5000); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.setParallelism(1);Properties props = new Properties(); props.setProperty("bootstrap.servers",KAFKA_BROKER); props.setProperty("zookeeper.connect", ZK_HOST); props.setProperty("group.id",GROUP_ID); props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); FlinkKafkaConsumer011<String> consumer = new FlinkKafkaConsumer011<>(TOPIC, new SimpleStringSchema(), props);consumer.setStartFromGroupOffsets();

注意:以下五種方式運行時優先級都比KafkaProperties中配置的auto.offset.reset優先級高。

方式一 : 指定topic, 指定partition的offset位置

Map<KafkaTopicPartition, Long> offsets = new HashedMap(); offsets.put(new KafkaTopicPartition("topic_name", 0), 11L); offsets.put(new KafkaTopicPartition("topic_name", 1), 22L); offsets.put(new KafkaTopicPartition("topic_name", 2), 33L); consumer.setStartFromSpecificOffsets(offsets);

Map<KafkaTopicPartition, Long> Long參數指定的offset位置

KafkaTopicPartition構造函數有兩個參數,第一個為topic名字,第二個為分區數.

  • 如果使用者需要讀取在提供的偏移量映射中沒有指定偏移量的分區,則它將回退到setStartFromGroupOffsets()該特定分區的默認組偏移行為。
  • 當作業從故障中自動恢復或使用保存點手動恢復時,這些起始位置配置方法不會影響起始位置。在恢復時,每個Kafka分區的起始位置由存儲在保存點或檢查點中的偏移量確定。

consumer.setStartFromSpecificOffsets(offsets);

方式二: 從topic中最初的數據開始消費

consumer.setStartFromEarliest();

方式三: 從指定的時間戳開始

consumer.setStartFromTimestamp(1559801580000l);

對于每個分區,時間戳大于或等于指定時間戳的記錄將用作起始位置。如果分區的最新記錄早于時間戳,則只會從最新記錄中讀取分區。在此模式下,Kafka中的已提交偏移將被忽略,不會用作起始位置。時間戳指的是kafka中消息自帶的時間戳。

方式四: 從最新的數據開始消費

consumer.setStartFromLatest();


方式五(同一默認)

參見: https://mp.weixin.qq.com/s?__biz=MzU5Mzk3MDA3Mw==&mid=2247483866&idx=2&sn=6a3b458caf5bebf0171f9fbd834b7517&chksm=fe09172cc97e9e3a590f5ea2978d078b1b46d94f86bd344173fa69c1d63790b09d2fe173bffb&token=1856795336&lang=zh_CN#rd

總結

以上是生活随笔為你收集整理的指定开始_Flink-Kafka指定offset的五种方式的全部內容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網站內容還不錯,歡迎將生活随笔推薦給好友。