Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -594,6 +594,7 @@
<allow pkg="org.apache.kafka.connect.runtime.distributed" />
<allow pkg="org.apache.kafka.connect.util" />
<allow pkg="org.apache.kafka.connect.converters" />
<allow pkg="org.apache.kafka.connect.json" />
<allow pkg="net.sourceforge.argparse4j" />
<!-- for tests -->
<allow pkg="org.apache.kafka.connect.integration" />
Expand Down Expand Up @@ -641,6 +642,7 @@
<allow pkg="org.apache.kafka.connect.util" />
<allow pkg="org.apache.kafka.common" />
<allow pkg="org.apache.kafka.connect.connector.policy" />
<allow pkg="org.apache.kafka.connect.json" />
</subpackage>

<subpackage name="storage">
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.common.utils.Exit;
import org.apache.kafka.connect.json.JsonConverter;
import org.apache.kafka.connect.json.JsonConverterConfig;
import org.apache.kafka.connect.mirror.rest.MirrorRestServer;
import org.apache.kafka.connect.runtime.Herder;
import org.apache.kafka.connect.runtime.isolation.Plugins;
Expand Down Expand Up @@ -51,6 +53,7 @@
import java.io.UnsupportedEncodingException;
import java.net.URLEncoder;
import java.nio.charset.StandardCharsets;
import java.util.Collections;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
Expand Down Expand Up @@ -269,7 +272,9 @@ private void addHerder(SourceAndTarget sourceAndTarget) {
Map<String, Object> adminProps = new HashMap<>(distributedConfig.originals());
ConnectUtils.addMetricsContextProperties(adminProps, distributedConfig, kafkaClusterId);
SharedTopicAdmin sharedAdmin = new SharedTopicAdmin(adminProps);
KafkaOffsetBackingStore offsetBackingStore = new KafkaOffsetBackingStore(sharedAdmin, () -> clientIdBase);
KafkaOffsetBackingStore offsetBackingStore = new KafkaOffsetBackingStore(sharedAdmin, () -> clientIdBase,
plugins.newInternalConverter(true, JsonConverter.class.getName(),
Collections.singletonMap(JsonConverterConfig.SCHEMAS_ENABLE_CONFIG, "false")));
offsetBackingStore.configure(distributedConfig);
ConnectorClientConfigOverridePolicy clientConfigOverridePolicy = new AllConnectorClientConfigOverridePolicy();
clientConfigOverridePolicy.configure(config.originals());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@

import org.apache.kafka.common.utils.Time;
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigOverridePolicy;
import org.apache.kafka.connect.json.JsonConverter;
import org.apache.kafka.connect.json.JsonConverterConfig;
import org.apache.kafka.connect.runtime.Herder;
import org.apache.kafka.connect.runtime.Worker;
import org.apache.kafka.connect.runtime.WorkerConfigTransformer;
Expand Down Expand Up @@ -77,7 +79,9 @@ protected Herder createHerder(DistributedConfig config, String workerId, Plugins
adminProps.put(CLIENT_ID_CONFIG, clientIdBase + "shared-admin");
SharedTopicAdmin sharedAdmin = new SharedTopicAdmin(adminProps);

KafkaOffsetBackingStore offsetBackingStore = new KafkaOffsetBackingStore(sharedAdmin, () -> clientIdBase);
KafkaOffsetBackingStore offsetBackingStore = new KafkaOffsetBackingStore(sharedAdmin, () -> clientIdBase,
plugins.newInternalConverter(true, JsonConverter.class.getName(),
Collections.singletonMap(JsonConverterConfig.SCHEMAS_ENABLE_CONFIG, "false")));
offsetBackingStore.configure(config);

Worker worker = new Worker(workerId, Time.SYSTEM, plugins, config, offsetBackingStore, connectorClientConfigOverridePolicy);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigOverridePolicy;
import org.apache.kafka.connect.json.JsonConverter;
import org.apache.kafka.connect.json.JsonConverterConfig;
import org.apache.kafka.connect.runtime.Connect;
import org.apache.kafka.connect.runtime.ConnectorConfig;
import org.apache.kafka.connect.runtime.Herder;
Expand All @@ -36,6 +38,7 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Collections;
import java.util.Map;

/**
Expand Down Expand Up @@ -89,7 +92,8 @@ protected Herder createHerder(StandaloneConfig config, String workerId, Plugins
ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy,
RestServer restServer, RestClient restClient) {

OffsetBackingStore offsetBackingStore = new FileOffsetBackingStore();
OffsetBackingStore offsetBackingStore = new FileOffsetBackingStore(plugins.newInternalConverter(
true, JsonConverter.class.getName(), Collections.singletonMap(JsonConverterConfig.SCHEMAS_ENABLE_CONFIG, "false")));
offsetBackingStore.configure(config);

Worker worker = new Worker(workerId, Time.SYSTEM, plugins, config, offsetBackingStore,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import org.apache.kafka.connect.runtime.rest.entities.ConfigKeyInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConfigValueInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorOffsets;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorStateInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorType;
import org.apache.kafka.connect.runtime.rest.errors.BadRequestException;
Expand Down Expand Up @@ -116,6 +117,7 @@ public abstract class AbstractHerder implements Herder, TaskStatus.Listener, Con
private final String kafkaClusterId;
protected final StatusBackingStore statusBackingStore;
protected final ConfigBackingStore configBackingStore;
protected ClusterConfigState configState;
private final ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy;
protected volatile boolean running = false;
private final ExecutorService connectorExecutor;
Expand All @@ -128,14 +130,28 @@ public AbstractHerder(Worker worker,
StatusBackingStore statusBackingStore,
ConfigBackingStore configBackingStore,
ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) {
this(worker, workerId, kafkaClusterId, statusBackingStore, configBackingStore, connectorClientConfigOverridePolicy,
Executors.newCachedThreadPool());
}

// Visible for testing
AbstractHerder(
Worker worker,
String workerId,
String kafkaClusterId,
StatusBackingStore statusBackingStore,
ConfigBackingStore configBackingStore,
ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy,
ExecutorService connectorExecutor) {
this.worker = worker;
this.worker.herder = this;
this.workerId = workerId;
this.kafkaClusterId = kafkaClusterId;
this.statusBackingStore = statusBackingStore;
this.configBackingStore = configBackingStore;
this.connectorClientConfigOverridePolicy = connectorClientConfigOverridePolicy;
this.connectorExecutor = Executors.newCachedThreadPool();
this.configState = ClusterConfigState.EMPTY;
this.connectorExecutor = connectorExecutor;
}

@Override
Expand Down Expand Up @@ -866,4 +882,20 @@ public List<ConfigKeyInfo> connectorPluginConfig(String pluginName) {
}
}

@Override
public void connectorOffsets(String connName, Callback<ConnectorOffsets> cb) {
Comment thread
C0urante marked this conversation as resolved.
log.debug("Submitting offset fetch request for connector: {}", connName);
Comment thread
yashmayya marked this conversation as resolved.
Outdated
connectorExecutor.submit(() -> {
Comment thread
yashmayya marked this conversation as resolved.
Outdated
try {
if (!configState.contains(connName)) {
Comment thread
yashmayya marked this conversation as resolved.
Outdated
cb.onCompletion(new NotFoundException("Connector " + connName + " not found"), null);
return;
}
ConnectorOffsets offsets = worker.connectorOffsets(connName, configState.connectorConfig(connName));
cb.onCompletion(null, offsets);
} catch (Throwable t) {
cb.onCompletion(t, null);
}
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.kafka.connect.runtime.rest.entities.ConfigInfos;
import org.apache.kafka.connect.runtime.rest.entities.ConfigKeyInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorOffsets;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorStateInfo;
import org.apache.kafka.connect.runtime.rest.entities.TaskInfo;
import org.apache.kafka.connect.storage.StatusBackingStore;
Expand Down Expand Up @@ -280,6 +281,13 @@ default void validateConnectorConfig(Map<String, String> connectorConfig, Callba
*/
List<ConfigKeyInfo> connectorPluginConfig(String pluginName);

/**
* Get the current offsets for a connector.
* @param connName the name of the connector whose offsets are to be retrieved
* @param cb callback to invoke upon completion
*/
void connectorOffsets(String connName, Callback<ConnectorOffsets> cb);

enum ConfigReloadAction {
NONE,
RESTART
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,14 +19,17 @@
import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.FenceProducersOptions;
import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsResult;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.KafkaFuture;
import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.MetricNameTemplate;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.common.config.ConfigValue;
import org.apache.kafka.common.config.provider.ConfigProvider;
Expand All @@ -42,6 +45,8 @@
import org.apache.kafka.connect.json.JsonConverterConfig;
import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup;
import org.apache.kafka.connect.runtime.isolation.LoaderSwap;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorOffset;
import org.apache.kafka.connect.runtime.rest.entities.ConnectorOffsets;
import org.apache.kafka.connect.runtime.rest.resources.ConnectResource;
import org.apache.kafka.connect.storage.ClusterConfigState;
import org.apache.kafka.connect.runtime.distributed.DistributedConfig;
Expand Down Expand Up @@ -88,6 +93,7 @@
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
Expand Down Expand Up @@ -1133,6 +1139,105 @@ public void setTargetState(String connName, TargetState state, Callback<TargetSt
}
}

/**
* Get the current offsets for a connector.
* @param connName the name of the connector whose offsets are to be retrieved
* @param connectorConfig the connector's configurations
* @return the connector's offsets
*/
public ConnectorOffsets connectorOffsets(String connName, Map<String, String> connectorConfig) {
Comment thread
C0urante marked this conversation as resolved.
Outdated
String connectorClassOrAlias = connectorConfig.get(ConnectorConfig.CONNECTOR_CLASS_CONFIG);
ClassLoader connectorLoader = plugins.connectorLoader(connectorClassOrAlias);
Connector connector;

try (LoaderSwap loaderSwap = plugins.withClassLoader(connectorLoader)) {
connector = plugins.newConnector(connectorClassOrAlias);
}
Comment thread
yashmayya marked this conversation as resolved.
Outdated

if (ConnectUtils.isSinkConnector(connector)) {
log.debug("Fetching offsets for sink connector: {}", connName);
return sinkConnectorOffsets(connName, connector, connectorConfig);
} else {
log.debug("Fetching offsets for source connector: {}", connName);
return sourceConnectorOffsets(connName, connector, connectorConfig);
}
}

/**
* Get the current consumer group offsets for a sink connector.
* @param connName the name of the sink connector whose offsets are to be retrieved
* @param connector the sink connector
* @param connectorConfig the sink connector's configurations
* @return the consumer group offsets for the sink connector
*/
private ConnectorOffsets sinkConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig) {
return sinkConnectorOffsets(connName, connector, connectorConfig, Admin::create);
}

// Visible for testing; allows us to mock out the Admin client for testing
ConnectorOffsets sinkConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig,
Function<Map<String, Object>, Admin> adminFactory) {
Map<String, Object> adminConfig = adminConfigs(
connName,
"connector-worker-adminclient-" + connName,
config,
new SinkConnectorConfig(plugins, connectorConfig),
connector.getClass(),
connectorClientConfigOverridePolicy,
kafkaClusterId,
ConnectorType.SOURCE);
Comment thread
yashmayya marked this conversation as resolved.
Outdated
String groupId = (String) baseConsumerConfigs(
connName, "connector-consumer-", config, new SinkConnectorConfig(plugins, connectorConfig),
connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SINK).get(ConsumerConfig.GROUP_ID_CONFIG);
Admin admin = adminFactory.apply(adminConfig);
Comment thread
C0urante marked this conversation as resolved.
try {
ListConsumerGroupOffsetsResult listConsumerGroupOffsetsResult = admin.listConsumerGroupOffsets(groupId);
Comment thread
yashmayya marked this conversation as resolved.
Outdated
try {
// Not using a timeout for the Future::get here because each offset get request is handled in its own thread in AbstractHerder
// and the REST API request timeout in HerderRequestHandler will ensure that the user request doesn't hang indefinitely
Comment thread
yashmayya marked this conversation as resolved.
Outdated
Map<TopicPartition, OffsetAndMetadata> offsets = listConsumerGroupOffsetsResult.all().get().get(groupId);
return SinkUtils.consumerGroupOffsetsToConnectorOffsets(offsets);
} catch (InterruptedException | ExecutionException e) {
throw new ConnectException("Failed to retrieve consumer group offsets for sink connector " + connName, e);
}
} finally {
Utils.closeQuietly(admin, "Offset fetch admin for sink connector " + connName);
}
}

/**
* Get the current offsets for a source connector.
* @param connName the name of the source connector whose offsets are to be retrieved
* @param connector the source connector
* @param connectorConfig the source connector's configurations
* @return the source connector's offsets
*/
private ConnectorOffsets sourceConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig) {
Comment thread
yashmayya marked this conversation as resolved.
Outdated
SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorConfig, config.topicCreationEnable());
ConnectorOffsetBackingStore offsetStore = config.exactlyOnceSourceEnabled()
? offsetStoreForExactlyOnceSourceConnector(sourceConfig, connName, connector)
: offsetStoreForRegularSourceConnector(sourceConfig, connName, connector);
CloseableOffsetStorageReader offsetReader = new OffsetStorageReaderImpl(offsetStore, connName, internalKeyConverter, internalValueConverter);
return sourceConnectorOffsets(connName, offsetStore, offsetReader);
}

// Visible for testing
ConnectorOffsets sourceConnectorOffsets(String connName, ConnectorOffsetBackingStore offsetStore,
CloseableOffsetStorageReader offsetReader) {
offsetStore.configure(config);
Comment thread
yashmayya marked this conversation as resolved.
Outdated
offsetStore.start();
try {
Set<Map<String, Object>> connectorPartitions = offsetStore.connectorPartitions(connName);
List<ConnectorOffset> connectorOffsets = offsetReader.offsets(connectorPartitions).entrySet().stream()
.map(entry -> new ConnectorOffset(entry.getKey(), entry.getValue()))
.collect(Collectors.toList());
return new ConnectorOffsets(connectorOffsets);
} finally {
Utils.closeQuietly(offsetReader, "Offset reader for connector " + connName);
offsetStore.stop();
Comment thread
yashmayya marked this conversation as resolved.
Outdated
}
}

ConnectorStatusMetricsGroup connectorStatusMetricsGroup() {
return connectorStatusMetricsGroup;
}
Expand Down Expand Up @@ -1426,7 +1531,7 @@ ConnectorOffsetBackingStore offsetStoreForRegularSourceConnector(

TopicAdmin admin = new TopicAdmin(adminOverrides);
KafkaOffsetBackingStore connectorStore =
KafkaOffsetBackingStore.forConnector(connectorSpecificOffsetsTopic, consumer, admin);
KafkaOffsetBackingStore.forConnector(connectorSpecificOffsetsTopic, consumer, admin, internalKeyConverter);

// If the connector's offsets topic is the same as the worker-global offsets topic, there's no need to construct
// an offset store that has a primary and a secondary store which both read from that same topic.
Expand Down Expand Up @@ -1483,7 +1588,7 @@ ConnectorOffsetBackingStore offsetStoreForExactlyOnceSourceConnector(

TopicAdmin admin = new TopicAdmin(adminOverrides);
KafkaOffsetBackingStore connectorStore =
KafkaOffsetBackingStore.forConnector(connectorSpecificOffsetsTopic, consumer, admin);
KafkaOffsetBackingStore.forConnector(connectorSpecificOffsetsTopic, consumer, admin, internalKeyConverter);

// If the connector's offsets topic is the same as the worker-global offsets topic, there's no need to construct
// an offset store that has a primary and a secondary store which both read from that same topic.
Expand Down Expand Up @@ -1532,7 +1637,7 @@ ConnectorOffsetBackingStore offsetStoreForRegularSourceTask(
KafkaConsumer<byte[], byte[]> consumer = new KafkaConsumer<>(consumerProps);

KafkaOffsetBackingStore connectorStore =
KafkaOffsetBackingStore.forTask(sourceConfig.offsetsTopic(), producer, consumer, topicAdmin);
KafkaOffsetBackingStore.forTask(sourceConfig.offsetsTopic(), producer, consumer, topicAdmin, internalKeyConverter);

// If the connector's offsets topic is the same as the worker-global offsets topic, there's no need to construct
// an offset store that has a primary and a secondary store which both read from that same topic.
Expand Down Expand Up @@ -1587,7 +1692,7 @@ ConnectorOffsetBackingStore offsetStoreForExactlyOnceSourceTask(
String connectorOffsetsTopic = Optional.ofNullable(sourceConfig.offsetsTopic()).orElse(config.offsetsTopic());

KafkaOffsetBackingStore connectorStore =
KafkaOffsetBackingStore.forTask(connectorOffsetsTopic, producer, consumer, topicAdmin);
KafkaOffsetBackingStore.forTask(connectorOffsetsTopic, producer, consumer, topicAdmin, internalKeyConverter);

// If the connector's offsets topic is the same as the worker-global offsets topic, there's no need to construct
// an offset store that has a primary and a secondary store which both read from that same topic.
Expand Down
Loading