如何在运行时为Kafka Topic动态添加前缀或后缀实现多环境隔离
Scala实现Kafka Topic运行时动态添加环境前缀/后缀方案
需求场景:当前运行环境为TEST1,原始Topic名称为BENEFIT,期望最终生成的Topic名称格式为
TEST1_BENEFIT,需要实现运行时为Kafka Topic动态添加前缀/后缀的能力。
实现思路
- 服务启动时一次性加载当前运行环境标识,避免每次发消息重复读取配置产生额外开销
- 统一封装Topic名称拼接逻辑,自动过滤空的前缀/后缀,避免出现多余下划线的非法Topic名
- 将拼接完成的最终Topic名传入
ProducerRecord,替代原始配置中的固定Topic值 - 补全原示例代码中缺失的序列化器配置、资源回收逻辑,避免基础运行错误
完整实现代码
import java.util.Properties import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} object Hello { // 全局加载运行环境标识,支持从配置文件/环境变量读取,示例值:TEST1、PROD、DEV private val runtimeEnv: String = configManager.getString("runtime.env") /** * 统一拼接Topic名称 * @param originTopic 原始配置的Topic名 * @param prefix 前缀,默认取当前运行环境 * @param suffix 后缀,默认空,可按需传入业务版本、分区标记等 * @return 拼接后的最终Topic名 */ private def wrapTopicName(originTopic: String, prefix: String = runtimeEnv, suffix: String = ""): String = { Seq(prefix, originTopic, suffix).filter(_.nonEmpty).mkString("_") } def Complete (jobId: BigInt, tableName: String, topic: String = configManager.getString("Kafka.Completion.Table.Topic")): Unit = { val kafkaServer = configManager.getString("Kafka.Server") val props = new Properties() props.put("bootstrap.servers", kafkaServer) // 补全原代码缺失的字符串序列化器配置 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") props.put("batch.size", "1") props.put("acks", "all") val producer = new KafkaProducer[String, String](props) // 生成带环境前缀的最终目标Topic val finalTargetTopic = wrapTopicName(topic) // 实际使用时替换为业务对应的key、value值 val message = new ProducerRecord[String, String](finalTargetTopic, key, value) producer.send(message) // 发送完成后回收Producer资源,生产环境建议使用全局单例Producer复用,无需每次创建关闭 producer.close() } }
注意事项
- 拼接逻辑已经做了空值过滤,如果本地开发不需要加前缀,只需要把
runtime.env配置为空字符串,就会直接使用原始Topic名,不会生成_BENEFIT这类格式错误的Topic - 如果需要添加后缀,比如业务版本标识,直接调用
wrapTopicName(originTopic = topic, suffix = "v2")即可生成TEST1_BENEFIT_v2格式的Topic - 生产环境不要每次发送消息都新建
KafkaProducer实例,建议初始化一个全局单例的Producer复用,大幅减少TCP连接建立、资源初始化的性能开销 - 上线前提前在Kafka集群创建好所有带环境前缀的Topic,或者开启Kafka自动创建Topic的权限,否则发送消息时会抛出Topic不存在的异常
内容的提问来源于stack exchange,提问作者dataeng
相关产品推荐
相关产品推荐

