Skip to content
Merged
Show file tree
Hide file tree
Changes from 10 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 @@ -866,4 +867,19 @@ public List<ConfigKeyInfo> connectorPluginConfig(String pluginName) {
}
}

@Override
public void connectorOffsets(String connName, Callback<ConnectorOffsets> cb) {
Comment thread
C0urante marked this conversation as resolved.
log.trace("Fetching offsets for connector: {}", connName);
ClusterConfigState configSnapshot = configBackingStore.snapshot();
try {
if (!configSnapshot.contains(connName)) {
cb.onCompletion(new NotFoundException("Connector " + connName + " not found"), null);
return;
}
// The worker asynchronously processes the request and completes the passed callback when done
worker.connectorOffsets(connName, configSnapshot.connectorConfig(connName), cb);
} 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,13 +19,15 @@
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.ListConsumerGroupOffsetsOptions;
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.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.KafkaFuture;
import org.apache.kafka.common.MetricNameTemplate;
import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.common.config.ConfigValue;
Expand All @@ -41,23 +43,25 @@
import org.apache.kafka.connect.json.JsonConverter;
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.resources.ConnectResource;
import org.apache.kafka.connect.storage.ClusterConfigState;
import org.apache.kafka.connect.runtime.distributed.DistributedConfig;
import org.apache.kafka.connect.runtime.errors.DeadLetterQueueReporter;
import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics;
import org.apache.kafka.connect.runtime.errors.ErrorReporter;
import org.apache.kafka.connect.runtime.errors.LogReporter;
import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator;
import org.apache.kafka.connect.runtime.errors.WorkerErrantRecordReporter;
import org.apache.kafka.connect.runtime.isolation.LoaderSwap;
import org.apache.kafka.connect.runtime.isolation.Plugins;
import org.apache.kafka.connect.runtime.isolation.Plugins.ClassLoaderUsage;
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.sink.SinkRecord;
import org.apache.kafka.connect.sink.SinkTask;
import org.apache.kafka.connect.source.SourceRecord;
import org.apache.kafka.connect.source.SourceTask;
import org.apache.kafka.connect.storage.CloseableOffsetStorageReader;
import org.apache.kafka.connect.storage.ClusterConfigState;
import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore;
import org.apache.kafka.connect.storage.Converter;
import org.apache.kafka.connect.storage.HeaderConverter;
Expand Down Expand Up @@ -1133,6 +1137,116 @@ public void setTargetState(String connName, TargetState state, Callback<TargetSt
}
}

/**
* Get the current offsets for a connector. This method is asynchronous and the passed callback is completed when the
* request finishes processing.
*
* @param connName the name of the connector whose offsets are to be retrieved
* @param connectorConfig the connector's configurations
* @param cb callback to invoke upon completion of the request
*/
public void connectorOffsets(String connName, Map<String, String> connectorConfig, Callback<ConnectorOffsets> cb) {
String connectorClassOrAlias = connectorConfig.get(ConnectorConfig.CONNECTOR_CLASS_CONFIG);
ClassLoader connectorLoader = plugins.connectorLoader(connectorClassOrAlias);

try (LoaderSwap loaderSwap = plugins.withClassLoader(connectorLoader)) {
Connector connector = plugins.newConnector(connectorClassOrAlias);
if (ConnectUtils.isSinkConnector(connector)) {
log.debug("Fetching offsets for sink connector: {}", connName);
sinkConnectorOffsets(connName, connector, connectorConfig, cb);
} else {
log.debug("Fetching offsets for source connector: {}", connName);
sourceConnectorOffsets(connName, connector, connectorConfig, cb);
}
}
}

/**
* 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
* @param cb callback to invoke upon completion of the request
*/
private void sinkConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig,
Callback<ConnectorOffsets> cb) {
sinkConnectorOffsets(connName, connector, connectorConfig, cb, Admin::create);
}

// Visible for testing; allows us to mock out the Admin client for testing
void sinkConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig,
Callback<ConnectorOffsets> cb, 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.SINK);
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 {
ListConsumerGroupOffsetsOptions listOffsetsOptions = new ListConsumerGroupOffsetsOptions()
.timeoutMs((int) ConnectResource.DEFAULT_REST_REQUEST_TIMEOUT_MS);
ListConsumerGroupOffsetsResult listConsumerGroupOffsetsResult = admin.listConsumerGroupOffsets(groupId, listOffsetsOptions);
listConsumerGroupOffsetsResult.partitionsToOffsetAndMetadata().whenComplete((result, error) -> {
if (error != null) {
log.error("Failed to retrieve consumer group offsets for sink connector {}", connName, error);
cb.onCompletion(new ConnectException("Failed to retrieve consumer group offsets for sink connector " + connName, error), null);
} else {
ConnectorOffsets offsets = SinkUtils.consumerGroupOffsetsToConnectorOffsets(result);
cb.onCompletion(null, offsets);
}
Utils.closeQuietly(admin, "Offset fetch admin for sink connector " + connName);
});
} catch (Throwable t) {
Utils.closeQuietly(admin, "Offset fetch admin for sink connector " + connName);
cb.onCompletion(new ConnectException("Failed to retrieve consumer group offsets for sink connector " + connName, t), null);
}
}

/**
* 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
* @param cb callback to invoke upon completion of the request
*/
private void sourceConnectorOffsets(String connName, Connector connector, Map<String, String> connectorConfig,
Callback<ConnectorOffsets> cb) {
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);
sourceConnectorOffsets(connName, offsetStore, offsetReader, cb);
}

// Visible for testing
void sourceConnectorOffsets(String connName, ConnectorOffsetBackingStore offsetStore,
CloseableOffsetStorageReader offsetReader, Callback<ConnectorOffsets> cb) {
executor.submit(() -> {
try {
offsetStore.configure(config);
offsetStore.start();
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());
cb.onCompletion(null, new ConnectorOffsets(connectorOffsets));
} catch (Throwable t) {
cb.onCompletion(t, null);
} finally {
Utils.closeQuietly(offsetReader, "Offset reader for connector " + connName);
Utils.closeQuietly(offsetStore::stop, "Offset store for connector " + connName);
}
});
}

ConnectorStatusMetricsGroup connectorStatusMetricsGroup() {
return connectorStatusMetricsGroup;
}
Expand Down Expand Up @@ -1426,7 +1540,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 +1597,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 +1646,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 +1701,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