【问题标题】:How to create simple kafka in kubernetes如何在 kubernetes 中创建简单的 kafka
【发布时间】:2021-07-27 14:10:24
【问题描述】:

我有一个任务是在 kubernetes 中使用 kafka 创建一个应用程序。但是当我将消费者连接到 kafka 时出现错误:

WARN 1 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-test-1, groupId=test] 连接到节点 -1 (kafka -service/10.99.233.131:9092) 无法建立。经纪人可能不可用。

WARN 1 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-test-1, groupId=test] 引导代理 kafka-service:9092 (id: -1 rack: null) 断开连接

这是我的 Kubernetes yaml 文件:

apiVersion: v1
kind: Service
metadata:
  name: kafka-service
  labels:
    name: kafka
spec:
  ports:
  - port: 9092
    targetPort: 9092
    protocol: TCP
  selector:
    name: kafka
---
apiVersion: v1
kind: Service
metadata:
  name: zookeeper-service
  labels:
    name: zookeeper
spec:
  ports:
  - name: client
    port: 2181
    protocol: TCP
  - name: follower
    port: 2888
    protocol: TCP
  - name: leader
    port: 3888
    protocol: TCP
  selector:
    name: zookeeper
  type: LoadBalancer
---
apiVersion: v1
kind: Pod
metadata:
  name: zookeeper
  labels:
    name: zookeeper 
spec:
  containers:
    - name: zookeeper
      image: zookeeper:3.7.0
---
apiVersion: v1
kind: Pod
metadata:
  name: kafka
  labels:
    name: kafka 
spec:
  containers:
    - name: kafka
      image: wurstmeister/kafka:2.13-2.6.0
      imagePullPolicy: "IfNotPresent"
      env:
        - name: KAFKA_ADVERTISED_PORT
          value: "666"
        - name: KAFKA_ADVERTISED_HOST_NAME
          value: localhost
        - name: KAFKA_ZOOKEEPER_CONNECT
          value: zookeeper-service:2181
        - name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
          value: INSIDE:PLAINTEXT
        - name: KAFKA_ADVERTISED_LISTENERS
          value: INSIDE://:666
        - name: KAFKA_LISTENERS
          value: INSIDE://:666
        - name: KAFKA_INTER_BROKER_LISTENER_NAME
          value: INSIDE
      ports:
        - containerPort: 9092

简单的java代码:

@Service
public class Consumer {
    @KafkaListener(topics = "new-topic",groupId = "test")
    public void consumeMessage(String message){
        System.out.println("************************");
        System.out.println(message);
        System.out.println("************************");
    }
}

如何在不扩展测试的情况下创建 kafka 代理?

【问题讨论】:

  • 从这里开始strimzi.io
  • 无论如何,您似乎误解了您设置的环境变量。 KAFKA_ADVERTISED_LISTENERS 优先于 KAFKA_ADVERTISED_PORT 和主机名。而且由于您已将侦听器端口设置为 666,因此您将无法在客户端中使用 9092 的容器端口

标签: java kubernetes apache-kafka


【解决方案1】:

这篇文章对我很有帮助:https://www.confluent.io/blog/kafka-listeners-explained/

我将 yaml 文件更改为:

apiVersion: v1
kind: Pod
metadata:
  name: kafka-service
  labels:
    name: kafka-service 
spec:
  containers:
    - name: kafka-service
      image: wurstmeister/kafka:2.13-2.6.0
      imagePullPolicy: "IfNotPresent"
      env:
        - name: KAFKA_ZOOKEEPER_CONNECT
          value: zookeeper-service:2181
        - name: KAFKA_BROKER_ID
          value: "1"

        - name: KAFKA_LISTENERS
          value: IN://:9092,OUT://:9093
        - name: KAFKA_ADVERTISED_LISTENERS
          value: IN://localhost:9092,OUT://kafka-service:9093
        - name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
          value: IN:PLAINTEXT,OUT:PLAINTEXT
        - name: KAFKA_INTER_BROKER_LISTENER_NAME
          value: IN

这意味着 9093 端口在 kubernetes 网络中打开,9092 端口从主机打开。 之后一切正常。

【讨论】:

    猜你喜欢
    • 2020-12-19
    • 2023-02-08
    • 1970-01-01
    • 2022-11-21
    • 2020-11-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-02-20
    相关资源
    最近更新 更多