You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 19:31:15