Kafka Sink Connector多实例变量异常原因求助
问题:Kafka Sink Connector中实例变量被覆盖的原因分析
问题背景
我是Java新手,开发Kafka Sink Connector时遇到一个问题,这应该是Java本身的问题而非Kafka专属。我部署了两个订阅不同Topic的连接器实例,但日志显示两个Topic的消息都使用了最后创建实例的influxMeasurement值。
相关代码
Connector类
package org.MySink.influxSink; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.connect.connector.Task; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.sink.SinkConnector; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class MySinkConnector extends SinkConnector { private static Logger log = LoggerFactory.getLogger(MySinkConnector.class); private Map<String, String> configProps; @Override public String version() { return VersionUtil.getVersion(); } @Override public void start(Map<String, String> map) { try { configProps = map; } catch(ConfigException e) { throw new ConnectException("Couldn't start InfluxSinkConnector due to configuration error", e); } } @Override public Class<? extends Task> taskClass() { return MySinkTask.class; } @Override public List<Map<String, String>> taskConfigs(int maxTasks) { List<Map<String, String>> taskConfigs = new ArrayList<>(); for(int i=0; i < maxTasks; i++) { taskConfigs.add(configProps); } return taskConfigs; } @Override public void stop() { return; } @Override public ConfigDef config() { return MySinkConnectorConfig.conf(); } }
Task类
package org.MySink.influxSink; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.sink.SinkRecord; import org.apache.kafka.connect.sink.SinkTask; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; import java.util.Collection; import java.util.HashSet; import java.util.Map; import java.util.Set; public class MySinkTask extends SinkTask { private static Logger log = LoggerFactory.getLogger(MySinkTask.class); private String influxMeasurement; private MySinkConnectorConfig config; private Map<String, String> configMap; @Override public String version() { return VersionUtil.getVersion(); } @Override public void start(Map<String, String> map) { config = new MySinkConnectorConfig(map); configMap = map; influxMeasurement = config.getInfluxMeasurement(); } @Override public void put(Collection<SinkRecord> collection) { if(collection.isEmpty()) { return; } final SinkRecord first = collection.iterator().next(); final int recordsCount = collection.size(); log.info(influxMeasurement + ": Received {} records. First record Kafka coordinates: ({}-{}-{}).", recordsCount, first.topic(), first.kafkaPartition(), first.kafkaOffset()); } @Override public void flush(Map<TopicPartition, OffsetAndMetadata> map) { } @Override public void stop() { //Close resources here. } }
异常现象
部署两个订阅不同Topic的连接器实例后,消息发布时日志显示:
kafka-connect | 2022-12-04T16:21:06.431482588Z [2022-12-04 16:21:06,431] INFO ActiveSessions: Received 1 records. First record Kafka coordinates: (TotalSessions-0-1134). (org.MySink.influxSink.MySinkTask) kafka-connect | 2022-12-04T16:21:06.431530001Z [2022-12-04 16:21:06,431] INFO ActiveSessions: Received 1 records. First record Kafka coordinates: (ActiveSessions-0-1122). (org.MySink.influxSink.MySinkTask)
可以看到,两个Topic的消息都显示使用最后创建实例的influxMeasurement值(都是ActiveSessions)。
解决方法
将Task类put方法中的日志语句从:
log.info(influxMeasurement + ": Received {} records. First record Kafka coordinates: ({}-{}-{}).", recordsCount, first.topic(), first.kafkaPartition(), first.kafkaOffset());
修改为:
log.info(this.influxMeasurement + ": Received {} records. First record Kafka coordinates: ({}-{}-{}).", recordsCount, first.topic(), first.kafkaPartition(), first.kafkaOffset());
修改后日志显示正常:
kafka-connect | 2022-12-04T16:21:06.431482588Z [2022-12-04 16:21:06,431] INFO TotalSessions: Received 1 records. First record Kafka coordinates: (TotalSessions-0-1134). (org.MySink.influxSink.MySinkTask) kafka-connect | 2022-12-04T16:21:06.431530001Z [2022-12-04 16:21:06,431] INFO ActiveSessions: Received 1 records. First record Kafka coordinates: (ActiveSessions-0-1122). (org.MySink.influxSink.MySinkTask)
原因分析
核心问题是你可能在MySinkConnectorConfig类中把influxMeasurement错误声明成了静态变量。
当你直接使用influxMeasurement时,如果存在同名的静态变量,Java会优先引用类级别的静态变量而非当前实例的成员变量。静态变量属于类本身,所有实例共享同一个存储值,所以第二个实例创建时会覆盖这个静态变量的值,导致所有实例都读取到最后设置的那个值。
加上this.关键字后,明确指定要引用当前实例的成员变量,避开了和静态变量的混淆,每个实例都会读取自己的influxMeasurement值,问题自然解决。
建议检查MySinkConnectorConfig类,把influxMeasurement改成实例变量,这样能从根源上避免这个问题,不用每次都手动加this.。
内容的提问来源于stack exchange,提问作者DiogoSilva14
相关产品推荐
相关产品推荐

