如何利用Kubernetes的ConfigMap功能存储Kafka Broker IP并供消费与生产者应用调用
使用Kubernetes ConfigMap管理Kafka Broker地址并在应用中调用
我来帮你梳理下从ConfigMap创建到生产者/消费者应用调用的全流程,都是实际可用的步骤和示例:
1. 创建存储Kafka Broker地址的ConfigMap
你可以通过两种方式创建ConfigMap:
方式一:通过YAML文件创建
创建一个名为kafka-brokers-config.yaml的文件,内容如下:
apiVersion: v1 kind: ConfigMap metadata: name: kafka-brokers-config namespace: your-namespace # 替换成你的应用所在命名空间 data: KAFKA_BOOTSTRAP_SERVERS: "192.168.1.10:9092,192.168.1.11:9092,192.168.1.12:9092"
然后执行命令创建:
kubectl apply -f kafka-brokers-config.yaml
方式二:直接用kubectl命令创建
如果不想写YAML,也可以直接通过命令行生成:
kubectl create configmap kafka-brokers-config --from-literal=KAFKA_BOOTSTRAP_SERVERS="192.168.1.10:9092,192.168.1.11:9092,192.168.1.12:9092" -n your-namespace
2. 在应用的Deployment中引用ConfigMap
接下来要让你的Kafka生产者/消费者应用能读取到这个ConfigMap的内容,有两种常用方式:
方式一:将ConfigMap内容注入为环境变量
修改你的应用Deployment YAML,在容器配置中添加环境变量引用:
apiVersion: apps/v1 kind: Deployment metadata: name: kafka-app-deployment namespace: your-namespace spec: replicas: 3 selector: matchLabels: app: kafka-app template: metadata: labels: app: kafka-app spec: containers: - name: kafka-app-container image: your-app-image:latest # 替换成你的应用镜像 env: - name: KAFKA_BOOTSTRAP_SERVERS valueFrom: configMapKeyRef: name: kafka-brokers-config key: KAFKA_BOOTSTRAP_SERVERS ports: - containerPort: 8080
方式二:将ConfigMap挂载为文件
如果你的应用习惯从配置文件读取,也可以把ConfigMap挂载成容器内的文件:
apiVersion: apps/v1 kind: Deployment metadata: name: kafka-app-deployment namespace: your-namespace spec: replicas: 3 selector: matchLabels: app: kafka-app template: metadata: labels: app: kafka-app spec: containers: - name: kafka-app-container image: your-app-image:latest volumeMounts: - name: kafka-config-volume mountPath: /app/config/kafka # 容器内的挂载路径 volumes: - name: kafka-config-volume configMap: name: kafka-brokers-config
挂载后,容器内/app/config/kafka/KAFKA_BOOTSTRAP_SERVERS文件里就会存储Broker地址的内容。
3. 在应用代码中读取配置
下面给出Java语言的示例(其他语言逻辑类似,都是读取环境变量或文件内容):
读取环境变量的示例(生产者)
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import java.util.Properties; public class KafkaProducerApp { public static void main(String[] args) { // 从环境变量读取Broker地址 String bootstrapServers = System.getenv("KAFKA_BOOTSTRAP_SERVERS"); Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); KafkaProducer<String, String> producer = new KafkaProducer<>(props); try { producer.send(new ProducerRecord<>("test-topic", "key1", "value1")).get(); System.out.println("消息发送成功"); } catch (Exception e) { e.printStackTrace(); } finally { producer.close(); } } }
读取挂载文件的示例(消费者)
import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.ConsumerRecord; import java.io.File; import java.io.FileNotFoundException; import java.util.Properties; import java.util.Arrays; import java.util.Scanner; public class KafkaConsumerApp { public static void main(String[] args) throws FileNotFoundException { // 读取挂载的配置文件内容 File configFile = new File("/app/config/kafka/KAFKA_BOOTSTRAP_SERVERS"); String bootstrapServers = new Scanner(configFile).useDelimiter("\\Z").next(); Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList("test-topic")); try { while (true) { ConsumerRecords<String, String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) { System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value()); } } } finally { consumer.close(); } } }
一些额外注意事项
- ConfigMap更新:如果后续Broker地址变化,只需要更新ConfigMap的内容,然后重启应用Pod即可生效;如果想要实现热重载,需要应用支持动态读取配置(比如用配置中心客户端,或者定时刷新文件内容)。
- 命名空间:确保ConfigMap和应用Deployment在同一个命名空间,或者在引用时指定正确的命名空间。
- K8s内部Kafka集群:如果你的Kafka是部署在K8s内部的,其实可以用Service的域名代替IP(比如
kafka-service:9092),这样Broker扩容缩容时不需要手动修改ConfigMap,不过你明确要求用IP,所以还是按你的需求来。
内容的提问来源于stack exchange,提问作者user1834664
相关产品推荐
相关产品推荐

