diff --git a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/consumer/CanalKafkaConsumer.java b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/consumer/CanalKafkaConsumer.java index 7aa8c28a..07be61d3 100644 --- a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/consumer/CanalKafkaConsumer.java +++ b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/consumer/CanalKafkaConsumer.java @@ -9,6 +9,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; +import com.alibaba.otter.canal.common.utils.PropertiesUtils; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -53,6 +54,8 @@ public class CanalKafkaConsumer implements CanalMsgConsumer { String k = (String) entry.getKey(); Object v = entry.getValue(); if (k.startsWith(PREFIX_KAFKA_CONFIG) && v != null) { + // check env config + v = PropertiesUtils.getProperty(properties, k); kafkaProperties.put(k.substring(PREFIX_KAFKA_CONFIG.length()), v); } } diff --git a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java index 0ba39cf1..3c30b96b 100644 --- a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java +++ b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java @@ -102,6 +102,8 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro String key = (String) entry.getKey(); Object value = entry.getValue(); if (key.startsWith(PREFIX_KAFKA_CONFIG) && value != null) { + // check env config + value = PropertiesUtils.getProperty(properties, key); key = key.substring(PREFIX_KAFKA_CONFIG.length()); kafkaProperties.put(key, value); }