You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Kubernetes的Apache Kafka本地开发配置与测试方法咨询

Kubernetes上Strimzi Kafka的自研应用对接方案与测试指南

一、核心连接配置解决端口转发/集群内接入问题

1. 集群内自研应用对接(推荐生产环境使用)

Strimzi部署的Kafka默认提供ClusterIP类型的bootstrap服务,同集群内应用直接通过服务名访问即可:

  • 同一Namespace下:使用{kafka-cluster-name}-kafka-bootstrap:9092作为Bootstrap Servers(替换{kafka-cluster-name}为你的Kafka集群名称)
  • 跨Namespace下:使用FQDN格式{kafka-cluster-name}-kafka-bootstrap.{kafka-namespace}.svc.cluster.local:9092
  • 无需额外端口转发,直接在应用配置中写入上述地址即可。

2. 本地开发机对接K8s内的Kafka

如果端口转发失败,大概率是Listener配置或转发端口错误,按以下步骤调整:

方法1:端口转发(快速测试用)

  • 转发Kafka内部Plain Listener的端口:
    kubectl port-forward svc/{kafka-cluster-name}-kafka-bootstrap 9092:9092
    
  • 本地应用配置Bootstrap Servers为localhost:9092
  • 注意:必须确保Kafka的Plain Listener(默认开启)允许内部访问,且未限制来源IP。

方法2:配置外部Listener(长期本地开发用)

修改Kafka CR文件,添加外部访问Listener(以NodePort为例,minikube/云环境通用):

apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
  name: {kafka-cluster-name}
spec:
  kafka:
    listeners:
      - name: plain
        port: 9092
        type: internal
        tls: false
      - name: external
        port: 9094
        type: NodePort
        tls: false # 测试阶段关闭TLS,生产环境建议开启
    # 保留原有的ZooKeeper、存储等配置
  zookeeper:
    replicas: 1
    storage:
      type: ephemeral
  • 应用配置Bootstrap Servers:
    • minikube环境:{minikube-ip}:{node-port}(通过minikube ip获取集群IP,kubectl get svc {kafka-cluster-name}-kafka-bootstrap查看NodePort)
    • 云环境:使用节点公网IP+NodePort,或改为type: LoadBalancer后用分配的外部IP+9094

3. TLS加密连接配置(生产环境必开)

如果开启了TLS Listener,需要将Strimzi自动生成的CA证书导入应用信任库:

  • 导出CA证书:
    kubectl get secret {kafka-cluster-name}-cluster-ca-cert -o jsonpath='{.data.ca\.crt}' | base64 -d > ca.crt
    
  • 应用配置中指定CA证书路径(以Java为例):
    props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL");
    props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, "/path/to/ca.crt");
    props.put(SslConfigs.SSL_TRUSTSTORE_TYPE_CONFIG, "PEM");
    

二、自研应用代码示例(Java/Kafka Clients)

生产者配置

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;

public class KafkaProducerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        // 集群内用服务名,本地用localhost:9092或外部IP+端口
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-cluster-kafka-bootstrap:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
            ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key1", "hello kafka");
            producer.send(record);
            System.out.println("消息发送成功");
        }
    }
}

消费者配置

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class KafkaConsumerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-cluster-kafka-bootstrap:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList("test-topic"));
            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> System.out.printf("收到消息:key=%s, value=%s%n", record.key(), record.value()));
            }
        }
    }
}

三、测试与业务逻辑验证步骤

1. 基础连通性测试

集群内测试

部署临时测试Pod,用Kafka命令行工具验证:

# 创建测试Pod
kubectl run -it --rm --image=confluentinc/cp-kafka:latest kafka-test-pod -- /bin/bash
# 在Pod内启动生产者
kafka-console-producer.sh --broker-list {kafka-cluster-name}-kafka-bootstrap:9092 --topic test-topic
# 新开终端启动消费者
kafka-console-consumer.sh --bootstrap-server {kafka-cluster-name}-kafka-bootstrap:9092 --topic test-topic --from-beginning

如果命令行能正常收发消息,说明Kafka集群本身无问题,再排查自研应用的配置。

本地测试

本地安装Kafka客户端工具,执行:

kafka-console-producer.sh --broker-list localhost:9092 --topic test-topic

确认命令行能连接后,再运行自研应用代码。

2. 日志排查

  • 查看Kafka Broker日志,定位连接错误:
    kubectl logs -f {kafka-cluster-name}-kafka-0
    
    重点搜索Connection refused、SSL handshake failed等关键词。
  • 查看自研应用Pod日志,检查是否有No route to host、Certificate not trusted等错误。

3. 业务逻辑验证

单元测试(脱离K8s集群)

使用嵌入式Kafka(如Spring Kafka的@EmbeddedKafka),直接测试生产消费的业务逻辑:

@SpringBootTest
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
public class KafkaBusinessTest {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @Test
    public void testMessageProcessing() throws InterruptedException {
        // 发送测试消息
        kafkaTemplate.send("test-topic", "test-key", "test-value");
        // 验证业务逻辑(如数据库写入、下游服务调用等)
        Thread.sleep(1000);
        // 断言结果
    }
}

集成测试(K8s环境)

  • 通过Strimzi的KafkaTopic CR创建测试主题:
    apiVersion: kafka.strimzi.io/v1beta2
    kind: KafkaTopic
    metadata:
      name: test-topic
      labels:
        strimzi.io/cluster: {kafka-cluster-name}
    spec:
      partitions: 3
      replicas: 1
      config:
        retention.ms: 7200000
    
  • 部署自研应用后,用脚本发送测试消息,验证应用是否正确消费并执行业务逻辑(如检查数据库数据、调用下游API等)。

内容的提问来源于stack exchange,提问作者GvidasP

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 16:57:20