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

Java分布式學習之Kafka消息隊列

 更新時間:2022年07月28日 11:33:13   作者:kaico2018  
Kafka是由Apache軟件基金會開發(fā)的一個開源流處理平臺,由Scala和Java編寫。Kafka是一種高吞吐量的分布式發(fā)布訂閱消息系統(tǒng),它可以處理消費者在網(wǎng)站中的所有動作流數(shù)據(jù)

介紹

Apache Kafka 是分布式發(fā)布-訂閱消息系統(tǒng),在 kafka官網(wǎng)上對 kafka 的定義:一個分布式發(fā)布-訂閱消息傳遞系統(tǒng)。 它最初由LinkedIn公司開發(fā),Linkedin于2010年貢獻給了Apache基金會并成為頂級開源項目。Kafka是一種快速、可擴展的、設計內(nèi)在就是分布式的,分區(qū)的和可復制的提交日志服務。

注意:Kafka并沒有遵循JMS規(guī)范(),它只提供了發(fā)布和訂閱通訊方式。

kafka中文官網(wǎng):http://kafka.apachecn.org/quickstart.html

Kafka核心相關名稱

  1. Broker:Kafka節(jié)點,一個Kafka節(jié)點就是一個broker,多個broker可以組成一個Kafka集群
  2. Topic:一類消息,消息存放的目錄即主題,例如page view日志、click日志等都可以以topic的形式存在,Kafka集群能夠同時負責多個topic的分發(fā)
  3. massage: Kafka中最基本的傳遞對象。
  4. Partition:topic物理上的分組,一個topic可以分為多個partition,每個partition是一個有序的隊列。Kafka里面實現(xiàn)分區(qū),一個broker就是表示一個區(qū)域。
  5. Segment:partition物理上由多個segment組成,每個Segment存著message信息
  6. Producer : 生產(chǎn)者,生產(chǎn)message發(fā)送到topic
  7. Consumer : 消費者,訂閱topic并消費message, consumer作為一個線程來消費
  8. Consumer Group:消費者組,一個Consumer Group包含多個consumer
  9. Offset:偏移量,理解為消息 partition 中消息的索引位置

主題和隊列的區(qū)別:

隊列是一個數(shù)據(jù)結構,遵循先進先出原則

kafka集群安裝

參考官方文檔:https://kafka.apachecn.org/quickstart.html

  • 每臺服務器上安裝jdk1.8環(huán)境
  • 安裝Zookeeper集群環(huán)境
  • 安裝kafka集群環(huán)境
  • 運行環(huán)境測試

安裝jdk環(huán)境和zookeeper這里不詳述了。

kafka為什么依賴于zookeeper:kafka會將mq信息存放到zookeeper上,為了使整個集群能夠方便擴展,采用zookeeper的事件通知相互感知。

kafka集群安裝步驟:

1、下載kafka的壓縮包,下載地址:https://kafka.apachecn.org/downloads.html

2、解壓安裝包

tar -zxvf kafka_2.11-1.0.0.tgz

3、修改kafka的配置文件 config/server.properties

配置文件修改內(nèi)容:

  • zookeeper連接地址:zookeeper.connect=192.168.1.19:2181
  • 監(jiān)聽的ip,修改為本機的iplisteners=PLAINTEXT://192.168.1.19:9092
  • kafka的brokerid,每臺broker的id都不一樣broker.id=0

4、依次啟動kafka

./kafka-server-start.sh -daemon config/server.properties

kafka使用

kafka文件存儲

topic是邏輯上的概念,而partition是物理上的概念,每個partition對應于一個log文件,該log文件中存儲的就是Producer生成的數(shù)據(jù)。Producer生成的數(shù)據(jù)會被不斷追加到該log文件末端,為防止log文件過大導致數(shù)據(jù)定位效率低下,Kafka采取了分片和索引機制,將每個partition分為多個segment,每個segment包括:“.index”文件、“.log”文件和.timeindex等文件。這些文件位于一個文件夾下,該文件夾的命名規(guī)則為:topic名稱+分區(qū)序號。

例如:執(zhí)行命令新建一個主題,分三個區(qū)存放放在三個broker中:

./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic kaico

  • 一個partition分為多個segment
  • .log 日志文件
  • .index 偏移量索引文件
  • .timeindex 時間戳索引文件
  • 其他文件(partition.metadata,leader-epoch-checkpoint)

Springboot整合kafka

maven依賴

 <dependencies>
        <!-- springBoot集成kafka -->
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
        </dependency>
        <!-- SpringBoot整合Web組件 -->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
    </dependencies>

yml配置

# kafka
spring:
  kafka:
    # kafka服務器地址(可以多個)
#    bootstrap-servers: 192.168.212.164:9092,192.168.212.167:9092,192.168.212.168:9092
    bootstrap-servers: www.kaicostudy.com:9092,www.kaicostudy.com:9093,www.kaicostudy.com:9094
    consumer:
      # 指定一個默認的組名
      group-id: kafkaGroup1
      # earliest:當各分區(qū)下有已提交的offset時,從提交的offset開始消費;無提交的offset時,從頭開始消費
      # latest:當各分區(qū)下有已提交的offset時,從提交的offset開始消費;無提交的offset時,消費新產(chǎn)生的該分區(qū)下的數(shù)據(jù)
      # none:topic各分區(qū)都存在已提交的offset時,從offset后開始消費;只要有一個分區(qū)不存在已提交的offset,則拋出異常
      auto-offset-reset: earliest
      # key/value的反序列化
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
    producer:
      # key/value的序列化
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      # 批量抓取
      batch-size: 65536
      # 緩存容量
      buffer-memory: 524288
      # 服務器地址
      bootstrap-servers: www.kaicostudy.com:9092,www.kaicostudy.com:9093,www.kaicostudy.com:9094

生產(chǎn)者

@RestController
public class KafkaController {
	/**
	 * 注入kafkaTemplate
	 */
	@Autowired
	private KafkaTemplate<String, String> kafkaTemplate;
	/**
	 * 發(fā)送消息的方法
	 *
	 * @param key
	 *            推送數(shù)據(jù)的key
	 * @param data
	 *            推送數(shù)據(jù)的data
	 */
	private void send(String key, String data) {
		// topic 名稱 key   data 消息數(shù)據(jù)
		kafkaTemplate.send("kaico", key, data);
	}
	// test 主題 1 my_test 3
	@RequestMapping("/kafka")
	public String testKafka() {
		int iMax = 6;
		for (int i = 1; i < iMax; i++) {
			send("key" + i, "data" + i);
		}
		return "success";
	}
}

消費者

@Component
public class TopicKaicoConsumer {
    /**
     * 消費者使用日志打印消息
     */
    @KafkaListener(topics = "kaico") //監(jiān)聽的主題
    public void receive(ConsumerRecord<?, ?> consumer) {
        System.out.println("topic名稱:" + consumer.topic() + ",key:" +
                consumer.key() + "," +
                "分區(qū)位置:" + consumer.partition()
                + ", 下標" + consumer.offset());
        //輸出key對應的value的值
        System.out.println(consumer.value());
    }
}

到此這篇關于Java分布式學習之Kafka消息隊列的文章就介紹到這了,更多相關Java Kafka內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • 解決springboot啟動成功,但訪問404的問題

    解決springboot啟動成功,但訪問404的問題

    這篇文章主要介紹了解決springboot啟動成功,但訪問404的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-07-07
  • 使用Spring Security和JWT實現(xiàn)安全認證機制

    使用Spring Security和JWT實現(xiàn)安全認證機制

    在現(xiàn)代 Web 應用中,安全認證和授權是保障數(shù)據(jù)安全和用戶隱私的核心機制,Spring Security 是 Spring 框架下專為安全設計的模塊,具有高度的可配置性和擴展性,而 JWT則是當前流行的認證解決方案,所以本文介紹了如何使用Spring Security和JWT實現(xiàn)安全認證機制
    2024-11-11
  • Springboot如何實現(xiàn)Web系統(tǒng)License授權認證

    Springboot如何實現(xiàn)Web系統(tǒng)License授權認證

    這篇文章主要介紹了Springboot如何實現(xiàn)Web系統(tǒng)License授權認證,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下
    2020-05-05
  • 詳解Java線程池的使用及工作原理

    詳解Java線程池的使用及工作原理

    在日常開發(fā)過程中總是以單線程的思維去編碼,沒有考慮到在多線程狀態(tài)下的運行狀況.由此引發(fā)的結果就是請求過多,應用無法響應.為了解決請求過多的問題,又衍生出了線程池的概念.本文記錄了Java中線程池的使用及工作原理,需要的朋友可以參考下
    2021-05-05
  • Springboot快速集成sse服務端推流(最新整理)

    Springboot快速集成sse服務端推流(最新整理)

    SSE?Server-Sent?Events是一種允許服務器向客戶端推送實時數(shù)據(jù)的技術,它建立在?HTTP?和簡單文本格式之上,提供了一種輕量級的服務器推送方式,通常也被稱為“事件流”(Event?Stream),這篇文章主要介紹了Springboot快速集成sse服務端推流(最新整理),需要的朋友可以參考下
    2024-02-02
  • java實現(xiàn)滑動驗證解鎖

    java實現(xiàn)滑動驗證解鎖

    這篇文章主要為大家詳細介紹了java實現(xiàn)滑動驗證解鎖,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-07-07
  • ruoyi微服務版本搭建運行方式

    ruoyi微服務版本搭建運行方式

    這篇文章主要介紹了ruoyi微服務版本搭建運行方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • Spring?@bean和@component注解區(qū)別

    Spring?@bean和@component注解區(qū)別

    本文主要介紹了Spring?@bean和@component注解區(qū)別,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-01-01
  • 深入解析Java的線程同步以及線程間通信

    深入解析Java的線程同步以及線程間通信

    這篇文章主要介紹了Java的線程同步以及線程間通信,多線程編程是Java學習中的重點和難點,需要的朋友可以參考下
    2015-09-09
  • 在Idea2020.1中使用gitee2020.1.0創(chuàng)建第一個代碼庫的實現(xiàn)

    在Idea2020.1中使用gitee2020.1.0創(chuàng)建第一個代碼庫的實現(xiàn)

    這篇文章主要介紹了在Idea2020.1中使用gitee2020.1.0創(chuàng)建第一個代碼庫的實現(xiàn),文中通過圖文示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2020-07-07

最新評論

临城县| 东明县| 大埔县| 阜新市| 双流县| 罗田县| 华亭县| 洞口县| 南木林县| 根河市| 肇庆市| 康平县| 邵阳县| 射洪县| 河东区| 肇东市| 九龙县| 乌拉特中旗| 西峡县| 沧州市| 苍山县| 班戈县| 台前县| 当涂县| 嘉义市| 星子县| 平和县| 湘潭市| 沐川县| 五家渠市| 通许县| 大化| 邯郸县| 山东| 霍邱县| 阳东县| 闽侯县| 蒲江县| 灵丘县| 商都县| 和林格尔县|