基于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
- minikube环境:
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-0Connection 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的
KafkaTopicCR创建测试主题: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
相关产品推荐
相关产品推荐

