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

Kafka調(diào)試技巧及心得分享

 更新時(shí)間:2026年02月12日 09:09:37   作者:程序員Forlan  
Kafka消費(fèi)組機(jī)制確保每個(gè)消息只被一個(gè)消費(fèi)者消費(fèi),新增機(jī)器時(shí)根據(jù)策略分配分區(qū),本地調(diào)試時(shí),修改代碼后可能需要重新測試和造消息,但可以通過配置參數(shù)控制消費(fèi)偏移量,線上環(huán)境調(diào)試可以通過接口拉取處理

基礎(chǔ)理念

對于我們有a,b,c,3臺機(jī)器,那么我們的消息會被消費(fèi)3次?

  • 可能會,也可能不會,這取決于你的配置和策略。
  • 消費(fèi)者組機(jī)制:Kafka 使用消費(fèi)組(Consumer Group)來確保每個(gè)消息只會被每個(gè)消費(fèi)者組中的一個(gè)消費(fèi)者消費(fèi)一次。
  • 分區(qū)分配:Kafka 主題可以分為多個(gè)分區(qū)(Partitions),每個(gè)分區(qū)只能由一個(gè)消費(fèi)者組中的一個(gè)消費(fèi)者消費(fèi)。

如果新增1臺機(jī)器,那么他的偏移量從0開始?

在kafka中,消費(fèi)組都會維護(hù)自己的偏移量(offset),以此來記錄消費(fèi)的消息位置,而當(dāng)新增1臺機(jī)器時(shí),會根據(jù)分區(qū)分配策略,比如:范圍分配、輪詢分配,這就有2種情況,可能加入舊分區(qū),也可能加入新分區(qū),具體可以配置下策略參數(shù)auto.offset.reset,對應(yīng)的值如下:

  • earliest: automatically reset the offset to the earliest offset
  • latest: automatically reset the offset to the latest offset
  • none: throw exception to the consumer if no previous offset is found for the consumer’s group
  • anything else: throw exception to the consumer.

本地調(diào)試

所以,當(dāng)我們在調(diào)試kafka消費(fèi)邏輯的時(shí)候,可能由于消費(fèi)邏輯寫的不對,改完代碼需要重新測,重新去造1條消息?還是你會怎么去做,了解了前面的原理,我們改消費(fèi)組是不可行的,他獲取的是最新的偏移量,無法實(shí)現(xiàn)復(fù)用之前造的某條數(shù)據(jù),特別是我們有不同邏輯,把每種類型的消息都重新推一波,這不僅麻煩,而且也容易出錯(cuò),是否可以直接復(fù)用之前的消息,準(zhǔn)確處理?

我們可以先了解下@KafkaListener里面的一些配置參數(shù),具體如下:

@KafkaListener(topicPartitions = {@TopicPartition(topic = "yourTopic", partitionOffsets = {@PartitionOffset(partition = "指定分區(qū)", initialOffset = "初始偏移量")})})
public void forlanConsumer(ConsumerRecord<String, String> record) {
	String messageStr = record.value();
	log.info(“測試觸達(dá)記錄消費(fèi):offset = {}, res = {}”,record.offset(), messageStr);
}

上面只是指定了我們從什么偏移量開始消費(fèi),如果要限制范圍,可以在代碼里面加限制

@KafkaListener(topicPartitions = {@TopicPartition(topic = "yourTopic", partitionOffsets = {@PartitionOffset(partition = "指定分區(qū)", initialOffset = "初始偏移量")})})
public void forlanConsumer(ConsumerRecord<String, String> record) {
	if (record.offset() > 結(jié)束偏移量) return;
	String messageStr = record.value();
	log.info(“測試觸達(dá)記錄消費(fèi):offset = {}, res = {}”,record.offset(), messageStr);
}

測試或正式環(huán)境調(diào)試

上面只適合本地場景,如果是線上環(huán)境,我們本地一般是沒有權(quán)限連接監(jiān)聽的,那么可以怎么做?其實(shí)也能做,只不過需要通過接口去拉取處理

@Autowired
private KafkaConfig kafkaConfig;

public void reconsumeMessage(String topic, int partition, long offset) {
		ConsumerFactory<Integer, String> consumerFactory = kafkaConfig.consumerFactory();
		Map<String, Object> configurationProperties = consumerFactory.getConfigurationProperties();

		Map<String, Object> customProps = new HashMap<>();
		customProps.put(ConsumerConfig.GROUP_ID_CONFIG, "reconsume-temp-group-" + UUID.randomUUID());
		customProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");// 防止 offset 不存在時(shí)報(bào)錯(cuò)
		customProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
		customProps.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "1");// 控制每次 poll 只拉取一條
		// 復(fù)用 kafkaConfig 中的基礎(chǔ)配置
		Map<String, Object> props = new HashMap<>(configurationProperties);
		props.putAll(customProps);

		try (Consumer<String, String> consumer = new KafkaConsumer<>(props)) {
			TopicPartition topicPartition = new TopicPartition(topic, partition);

			// 分配分區(qū)并定位偏移量
			consumer.assign(Collections.singletonList(topicPartition));
			consumer.seek(topicPartition, offset);

			// 拉取消息(設(shè)置超時(shí)時(shí)間)
			ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(5));

			for (ConsumerRecord<String, String> record : records) {
				if (record.offset() == offset) {
					if (Objects.equals(topic, KafkaTopic.Forlan_MESSAGE_NOTIFY)) {
						// 調(diào)用@KafkaListener的方法執(zhí)行邏輯
						forlanConsumer.userConsumer(record);
					}
					break;
				}
			}
		}
	}

項(xiàng)目中如果沒有配置kafkaConfig,也可以自定義一個(gè),只要能拿到連接就行

private Map<String, Object> consumerConfigs() {
    Map<String, Object> props = new HashMap();
    props.put("bootstrap.servers", this.bootstrapServers);
    props.put("group.id", this.groupid);
    props.put("enable.auto.commit", this.autoCommit);
    props.put("auto.commit.interval.ms", this.interval);
    props.put("session.timeout.ms", this.timeout);
    props.put("key.deserializer", this.keyDeserializer);
    props.put("value.deserializer", this.valueDeserializer);
    props.put("auto.offset.reset", this.offsetReset);
    props.put("max.poll.records", this.maxPollRecords);
    return props;
}

總結(jié)

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

相關(guān)文章

  • Java打包可執(zhí)行JAR文件的三種方式詳解

    Java打包可執(zhí)行JAR文件的三種方式詳解

    很多?Java?開發(fā)者都遇到過這樣的問題,一打包就報(bào)錯(cuò)?ClassNotFoundException,其實(shí)這就是JAR?包沒打好,下面我們就來看看Maven打可執(zhí)行?JAR?的三種詳細(xì)方法吧
    2025-12-12
  • Windows系統(tǒng)下修改jar包中的依賴實(shí)現(xiàn)過程

    Windows系統(tǒng)下修改jar包中的依賴實(shí)現(xiàn)過程

    本文介紹了如何在Windows下修改和升級Java項(xiàng)目中的JAR包依賴,首先,通過解壓原JAR包并替換指定的依賴JAR文件來實(shí)現(xiàn)升級,然后,使用命令重新打包JAR文件,整個(gè)過程包括解壓、替換和打包三個(gè)步驟,確保項(xiàng)目能夠正確引用更新的依賴
    2025-12-12
  • Java多線程之多線程異常捕捉

    Java多線程之多線程異常捕捉

    在java多線程程序中,所有線程都不允許拋出未捕獲的checked exception,也就是說各個(gè)線程需要自己把自己的checked exception處理掉,通過此篇文章給大家分享Java多線程之多線程異常捕捉,需要的朋友可以參考下
    2015-08-08
  • 線程池之newCachedThreadPool可緩存線程池的實(shí)例

    線程池之newCachedThreadPool可緩存線程池的實(shí)例

    這篇文章主要介紹了線程池之newCachedThreadPool可緩存線程池的實(shí)例,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • 詳解Java中的防抖和節(jié)流

    詳解Java中的防抖和節(jié)流

    防抖是將多次執(zhí)行變?yōu)橹付〞r(shí)間內(nèi)不在觸發(fā)之后,執(zhí)行一次。節(jié)流是將多次執(zhí)行變?yōu)橹付〞r(shí)間不論觸發(fā)多少次,時(shí)間一到就執(zhí)行一次。這篇文章來和大家聊聊Java中的防抖和節(jié)流,感興趣的可以了解一下
    2022-08-08
  • java void方法單測斷言的三種實(shí)現(xiàn)示例

    java void方法單測斷言的三種實(shí)現(xiàn)示例

    Java中對void方法單元測試可通過驗(yàn)證狀態(tài)變化、行為交互及異常拋出實(shí)現(xiàn),推薦使用AssertJ和Mockito,文中通過示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2025-07-07
  • Spring中DAO被循環(huán)調(diào)用的時(shí)候數(shù)據(jù)不實(shí)時(shí)更新的解決方法

    Spring中DAO被循環(huán)調(diào)用的時(shí)候數(shù)據(jù)不實(shí)時(shí)更新的解決方法

    這篇文章主要介紹了Spring中DAO被循環(huán)調(diào)用的時(shí)候數(shù)據(jù)不實(shí)時(shí)更新的解決方法,需要的朋友可以參考下
    2014-08-08
  • nodejs與JAVA應(yīng)對高并發(fā)的對比方式

    nodejs與JAVA應(yīng)對高并發(fā)的對比方式

    這篇文章主要介紹了nodejs與JAVA應(yīng)對高并發(fā)的對比方式,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-08-08
  • 詳解基于spring多數(shù)據(jù)源動態(tài)調(diào)用及其事務(wù)處理

    詳解基于spring多數(shù)據(jù)源動態(tài)調(diào)用及其事務(wù)處理

    本篇文章主要介紹了基于spring多數(shù)據(jù)源動態(tài)調(diào)用及其事務(wù)處理 ,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-06-06
  • java中BigDecimal類的構(gòu)造詳解及使用

    java中BigDecimal類的構(gòu)造詳解及使用

    這篇文章主要介紹了java中BigDecimal類的構(gòu)造詳解及使用,Java在java.math包中提供的API類BigDecimal,用來對超過16位有效位的數(shù)進(jìn)行精確的運(yùn)算,需要的朋友可以參考下
    2023-07-07

最新評論

府谷县| 普兰县| 株洲县| 崇礼县| 凤庆县| 农安县| 曲周县| 耒阳市| 夏津县| 嘉定区| 海晏县| 无棣县| 霍州市| 堆龙德庆县| 通榆县| 晴隆县| 台南市| 岱山县| 玛曲县| 图们市| 宜君县| 乡宁县| 新昌县| 盐城市| 藁城市| 滕州市| 津市市| 保亭| 如皋市| 射洪县| 绵竹市| 申扎县| 丰镇市| 富川| 河南省| 竹溪县| 乐陵市| 青田县| 扎赉特旗| 海阳市| 延寿县|