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

docker安裝單機版kafka并使用的詳細步驟

 更新時間:2025年06月11日 10:07:18   作者:小咖張  
這篇文章主要為大家詳細介紹了docker安裝單機版kafka并使用的詳細步驟,文中的示例代碼講解詳細,感興趣的小伙伴可以跟隨小編一起學習一下

一、docker-compose.yml

version: '3'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:latest
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - "2181:2181"
    volumes:
      - ./zookeeper-data:/var/lib/zookeeper/data
      - ./zookeeper-log:/var/lib/zookeeper/log

  kafka:
    image: confluentinc/cp-kafka:latest
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://自己的ip:9092
      KAFKA_LISTENERS: PLAINTEXT://:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1

      #新版使用CFG

      KAFKA_CFG_PROCESS_ROLES: broker
      KAFKA_CFG_CONTROLLER_LISTENER_NAMES: PLAINTEXT
      KAFKA_CFG_LISTENERS: PLAINTEXT://:9092
      KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://自己的ip:9092
      KAFKA_CFG_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_CFG_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_CFG_DEFAULT_REPLICATION_FACTOR: 1
      KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_CFG_NUM_PARTITIONS: 1
      KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "true"
    volumes:
      - ./kafka-data:/var/lib/kafka/data

#kafka可視化界面

kafka-manager:
     image: hlebalbau/kafka-manager:latest
     ports:
        - "9000:9000"
     environment:
        ZK_HOSTS: "zookeeper:2181"
        APPLICATION_SECRET: "random-secret-key"
        KAFKA_MANAGER_LOG_LEVEL: "INFO"
     depends_on:
        - zookeeper
        - kafka

二、創(chuàng)建權(quán)限

sudo chown -R 1000:1000 ./zookeeper-data
sudo chown -R 1000:1000 ./zookeeper-log
sudo chown -R 1000:1000 ./kafka-data

三、啟動容器

docker-compose up -d

四、測試功能

1. 進入Kafka容器:

docker exec -it kafka bash

2. 創(chuàng)建一個測試主題:

kafka-topics --create --topic order1-topic --bootstrap-server 上面配置的ip:9092 --replication-factor 1 --partitions 1

3: 啟動一個生產(chǎn)者:

kafka-console-producer --topic order1-topic --bootstrap-server 上面配置的ip:9092

4. 在另一個終端窗口,啟動一個消費者:

docker exec -it kafka kafka-console-consumer --topic order1-topic --bootstrap-server  上面配置的ip:9092 --from-beginning

5:測試成功

生產(chǎn)者:

消費者:

五、kafka 可視化界面使用

打開界面進行添加Cluster,添加成功如下圖所示:

點擊進入可以進行Topics和Brokers的管理

六、項目中進行配置和使用

1:pom設(shè)置,因為我后面要用Flink,所以kafka.就直接引用flink-connector-kafka

<!--實時訂單處理-->
<!-- Flink 核心依賴 -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-java</artifactId>
    <version>${flink.version}</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-streaming-java</artifactId>
    <version>${flink.version}</version>
</dependency>

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients</artifactId>
    <version>${flink.version}</version>
</dependency>

<!-- Kafka Source Connector -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka</artifactId>
    <version>${flink.version}</version>
</dependency>

<!-- Redis Sink 客戶端 -->
<dependency>
    <groupId>redis.clients</groupId>
    <artifactId>jedis</artifactId>
    <version>3.9.0</version>
</dependency>

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

<!--實時訂單處理-->

application-dev.yml: 設(shè)置后,Consumer可以直接使用@KafkaListener監(jiān)聽消息

spring:
  kafka:
    listener:
      missing-topics-fatal: false
    bootstrap-servers: 自己的ip:9092
    consumer:
      auto-offset-reset: earliest
      enable-auto-commit: false
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
    properties:
      session.timeout.ms: 15000
      heartbeat.interval.ms: 5000
      max.poll.interval.ms: 300000
      metadata.max.age.ms: 3000
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer

2:多 Topic + 多 Group 實現(xiàn)方式:

功能實現(xiàn)方式
多 Topic多個 @KafkaListener 方法
多 Group每個監(jiān)聽器指定不同的 groupId
負載均衡多實例使用相同 groupId
廣播模式多實例使用不同 groupId
動態(tài)配置使用 application.yml + @Value 注入
自定義容器工廠使用 ConcurrentKafkaListenerContainerFactory

3:示例代碼

KafkaConfig:

package com.zbkj.front.config.kafka;
 
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.*;
 
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
 
@Configuration
@EnableKafka
public class KafkaConfig {
 
    // ========== 公共方法 ==========
    public Map<String, Object> commonConsumerProps(String groupId) {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "自己的ip:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        return props;
    }
 
    // ========== 消費者組 1 - order-group ==========
    @Bean
    public ConsumerFactory<String, String> orderConsumerFactory() {
        return new DefaultKafkaConsumerFactory<>(commonConsumerProps("order-group"));
    }
 
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> orderKafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(orderConsumerFactory());
        factory.setConcurrency(1); // 可根據(jù)分區(qū)數(shù)設(shè)置并發(fā)度
        return factory;
    }
 
    // ========== 消費者組 2 - payment-group ==========
    @Bean
    public ConsumerFactory<String, String> paymentConsumerFactory() {
        return new DefaultKafkaConsumerFactory<>(commonConsumerProps("payment-group"));
    }
 
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> paymentKafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(paymentConsumerFactory());
        factory.setConcurrency(1);
        return factory;
    }
 
    // ========== Kafka Producer Bean ==========
    @Bean
    public KafkaProducer<String, String> orderKafkaProducer() {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "自己的ip:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "5");
 
        return new KafkaProducer<>(props);
    }
}

OrderConsumer

package com.zbkj.front.config.kafka;
 
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
 
@Service
@Slf4j
public class OrderConsumer {
 
    @KafkaListener(topics = "orderCreate-topic", containerFactory = "orderKafkaListenerContainerFactory")
    public void consumeOrder(ConsumerRecord<String, String> record) {
        log.info("[訂單服務(wù)] 收到消息 topic={}, offset={}, key={}, value={}",
                record.topic(), record.offset(), record.key(), record.value());
    }
}

PaymentConsumer

package com.zbkj.front.config.kafka;
 
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
 
@Service
@Slf4j
public class PaymentConsumer {
 
    @KafkaListener(topics = "orderPayment-topic", containerFactory = "paymentKafkaListenerContainerFactory")
    public void consumePayment(ConsumerRecord<String, String> record) {
        log.info("[支付服務(wù)] 收到消息 topic={}, offset={}, key={}, value={}",
                record.topic(), record.offset(), record.key(), record.value());
    }
}

testController

package com.zbkj.front.controller;
 
import com.fasterxml.jackson.databind.ObjectMapper;
import com.zbkj.common.model.order.Order;
import com.zbkj.common.result.CommonResult;
import com.zbkj.front.event.OrderEvent;
import com.zbkj.front.event.OrderStatusEum;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.RestController;
 
import java.math.BigDecimal;
import java.time.LocalDateTime;
 
@Slf4j
@RestController
@RequestMapping("api/front/test")
@Api(tags = "測試")
public class testController {
 
    private final ObjectMapper objectMapper = new ObjectMapper();
    @Autowired
    @Qualifier("orderKafkaProducer")
    private KafkaProducer<String, String> orderKafkaProducer;
 
 
    @ApiOperation(value = "測試")
    @RequestMapping(value = "/test/orderCreate", method = RequestMethod.GET)
    public CommonResult<String> orderTopic(@Validated String orderNo) {
        Order order= new Order();
        order.setOrderNo(orderNo);
        order.setPayPrice(new BigDecimal(100));
        // 添加到訂單創(chuàng)建的地方
        sendOrderCreatedEvent(order,"orderCreate-topic");
      return   CommonResult.success();
    }
 
    @ApiOperation(value = "測試")
    @RequestMapping(value = "/test/orderPayment", method = RequestMethod.GET)
    public CommonResult<String> orderPayment(@Validated String orderNo) {
        Order order= new Order();
        order.setOrderNo(orderNo);
        order.setPayPrice(new BigDecimal(100));
        // 添加到訂單創(chuàng)建的地方
        sendOrderCreatedEvent(order,"orderPayment-topic");
        return   CommonResult.success();
    }
 
    public void sendOrderCreatedEvent(Order order,String topic) {
        try {
            // 構(gòu)建訂單事件對象
            OrderEvent event = new OrderEvent(
                    String.valueOf(order.getOrderNo()),
                    order.getPayPrice(),
                    LocalDateTime.now().toString(),
                    OrderStatusEum.CREATED
            );
 
            // 序列化為 JSON
            String json = objectMapper.writeValueAsString(event);
 
            // 創(chuàng)建 Kafka 消息記錄(使用訂單號作為 key)
            ProducerRecord<String, String> record = new ProducerRecord<>(topic, order.getOrderNo(), json);
 
            // 發(fā)送消息
            orderKafkaProducer.send(record, (metadata, exception) -> {
                if (exception != null) {
                    log.error("Kafka消息發(fā)送失敗 topic={}, key={}, error={}",
                            record.topic(), record.key(), exception.getMessage());
                } else {
                    log.info("Kafka消息發(fā)送成功 topic={}, key={}, partition={}, offset={}",
                            metadata.topic(), record.key(), metadata.partition(), metadata.offset());
                }
            });
 
        } catch (Exception e) {
            log.error("發(fā)送Kafka訂單創(chuàng)建事件異常 orderNo={}", order.getOrderNo(), e);
        }
    }
 
}

執(zhí)行完后:

2025-05-28 15:56:53.768 [kafka-producer-network-thread | producer-1] INFO  com.zbkj.front.controller.testController - Kafka消息發(fā)送成功 topic=orderCreate-topic, key=PT672174822506216634598, partition=0, offset=0
2025-05-28 15:56:53.772 [org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1] INFO  com.zbkj.front.config.kafka.OrderConsumer - [訂單服務(wù)] 收到消息 topic=orderCreate-topic, offset=0, key=PT672174822506216634598, value={"orderId":"PT672174822506216634598","merId":null,"amount":100,"timestamp":"2025-05-28T15:56:49.726","status":"CREATED"}

2025-05-28 15:57:22.687 [kafka-producer-network-thread | producer-1] INFO  com.zbkj.front.controller.testController - Kafka消息發(fā)送成功 topic=orderPayment-topic, key=SH377174799246359695171, partition=0, offset=0
2025-05-28 15:57:22.688 [org.springframework.kafka.KafkaListenerEndpointContainer#1-0-C-1] INFO  com.zbkj.front.config.kafka.PaymentConsumer - [支付服務(wù)] 收到消息 topic=orderPayment-topic, offset=0, key=SH377174799246359695171, value={"orderId":"SH377174799246359695171","merId":null,"amount":100,"timestamp":"2025-05-28T15:57:20.754","status":"CREATED"}

以上就是docker安裝單機版kafka并使用的詳細步驟的詳細內(nèi)容,更多關(guān)于docker安裝kafka的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Docker?Compose資源配額管理實現(xiàn)限制容器CPU與內(nèi)存使用

    Docker?Compose資源配額管理實現(xiàn)限制容器CPU與內(nèi)存使用

    在容器化部署中,資源競爭是最常見的穩(wěn)定性隱患,當一個服務(wù)異常占用CPU或內(nèi)存時,可能導致整個應(yīng)用集群異常,本文將系統(tǒng)講解Docker?Compose的資源配額管理方案,感興趣的可以了解一下
    2026-04-04
  • docker?compose運行微服務(wù)項目的方法

    docker?compose運行微服務(wù)項目的方法

    這篇文章主要介紹了docker?compose運行微服務(wù)項目?,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-08-08
  • FastAPI 部署在Docker的詳細過程

    FastAPI 部署在Docker的詳細過程

    這篇文章主要介紹了FastAPI 部署在 Docker的詳細過程,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-10-10
  • Docker如何搭建私有倉庫

    Docker如何搭建私有倉庫

    文章介紹了如何搭建私有倉庫并使用Docker進行鏡像的管理和推送,首先,搭建私有倉庫并配置非HTTPS訪問(適用于Ubuntu、Debian和CentOS),然后,使用Docker命令標記、推送和拉取鏡像,最后,通過curl命令查看倉庫中的鏡像列表
    2025-03-03
  • Docker安裝LNMP環(huán)境的詳細過程(可部署TP項目)

    Docker安裝LNMP環(huán)境的詳細過程(可部署TP項目)

    這篇文章主要介紹了Docker安裝LNMP環(huán)境的詳細過程(可部署TP項目),主要包括安裝docker,安裝nginx,安裝php的命令詳解,需要的朋友可以參考下
    2022-06-06
  • 使用docker制作分布式lnmp 鏡像

    使用docker制作分布式lnmp 鏡像

    最近在學習docker相關(guān)知識,順便把docker制作分布式lnmp 鏡像的過程分享給大家,包括Nginx配置文件和PHP文件的修改代碼也一并給出,感興趣的朋友跟隨小編一起看看吧
    2021-06-06
  • 在OpenWrt設(shè)備上搭建Docker環(huán)境的完整方案

    在OpenWrt設(shè)備上搭建Docker環(huán)境的完整方案

    OpenWrt 作為一個高度可定制的嵌入式 Linux 發(fā)行版,其模塊化設(shè)計為 Docker 容器化部署提供了可能性,將 Docker 引入 OpenWrt 環(huán)境,能夠帶來許多優(yōu)勢,所以本文給大家分享了在 OpenWrt 設(shè)備上搭建 Docker 環(huán)境的完整方案,需要的朋友可以參考下
    2025-04-04
  • docker鏡像管理命令詳解

    docker鏡像管理命令詳解

    這篇文章主要介紹了docker鏡像管理命令,我們也可以使用命令來搜索鏡像,比如我們需要一個tomcat的鏡像來作為我們的web服務(wù),我們可以通過 docker search 命令搜索tomcat來尋找適合我們的鏡像,本文給大家介紹的非常詳細,需要的朋友參考下吧
    2022-07-07
  • 減少Docker鏡像大小的10個優(yōu)化技巧

    減少Docker鏡像大小的10個優(yōu)化技巧

    當使用Docker時,鏡像大小是一個很大的問題,下面這篇文章主要給大家介紹了關(guān)于減少Docker鏡像大小的10個優(yōu)化技巧,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2024-01-01
  • Docker更換鏡像源詳細代碼教程

    Docker更換鏡像源詳細代碼教程

    Docker是一個開源的應(yīng)用容器引擎,使用Go語言編寫,允許開發(fā)者將應(yīng)用及依賴打包到輕量級容器中,可在不同Linux系統(tǒng)間移植,這篇文章主要給大家介紹了關(guān)于Docker更換鏡像源的相關(guān)資料,需要的朋友可以參考下
    2024-08-08

最新評論

浠水县| 丰宁| 吴川市| 武川县| 娄底市| 库尔勒市| 平顺县| 通河县| 新绛县| 广宗县| 曲沃县| 贵德县| 那曲县| 湖北省| 宝应县| 双江| 宝鸡市| 宜州市| 吉安县| 玉环县| 吐鲁番市| 蓬溪县| 盱眙县| 海淀区| 广东省| 鄂托克前旗| 报价| 新宁县| 寿阳县| 鹤壁市| 兰坪| 竹山县| 宕昌县| 阜宁县| 丹江口市| 香港| 崇阳县| 古浪县| 阿鲁科尔沁旗| 城步| 右玉县|