Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
2d5ea1c
Add IsolatedPlugin
gharris1727 Jan 27, 2023
768bb22
Add Isolated*Connector classes
gharris1727 Jan 27, 2023
367f2cb
Return IsolatedConnector from Plugins
gharris1727 Jan 27, 2023
fb07a9b
Use Isolated*Connector in Herder and Worker classes
gharris1727 Feb 1, 2023
e35c06f
fixup: remove unnecessary catch clauses
gharris1727 Feb 1, 2023
72cf9c7
fixup: revert config transformer error handling
gharris1727 Feb 2, 2023
84acab0
fixup: review comments for IsolatedPlugin part 1
gharris1727 Mar 1, 2023
3780596
fixup: Review comments for Isolated*Connector part 1
gharris1727 Mar 1, 2023
05ce444
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Mar 2, 2023
8035077
fixup: remove checkstyle suppression
gharris1727 Mar 2, 2023
bb99bad
fixup: merge conflict
gharris1727 Mar 2, 2023
493bd0a
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Jun 8, 2023
ee93a63
fixup: merge conflict
gharris1727 Jun 8, 2023
5906ed6
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Oct 13, 2023
8cd9d11
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Jan 3, 2024
561c174
Merge StandaloneHerderTest
gharris1727 Jan 3, 2024
2abb38c
fixup: fix this-escape warning on jdk21
gharris1727 Jan 3, 2024
2bf6ea5
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Jan 12, 2024
8df55c0
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 Feb 12, 2024
15f3e46
Merge remote-tracking branch 'upstream/trunk' into kafka-14670-wrap-c…
gharris1727 May 7, 2024
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
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,11 @@
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigRequest;
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.errors.NotFoundException;
import org.apache.kafka.connect.runtime.isolation.IsolatedConnector;
import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector;
import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector;
import org.apache.kafka.connect.runtime.isolation.LoaderSwap;
import org.apache.kafka.connect.runtime.isolation.PluginType;
import org.apache.kafka.connect.runtime.isolation.Plugins;
import org.apache.kafka.connect.runtime.rest.entities.ActiveTopicsInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConfigInfo;
Expand Down Expand Up @@ -137,7 +141,7 @@ public abstract class AbstractHerder implements Herder, TaskStatus.Listener, Con
private final Time time;
protected final Loggers loggers;

private final ConcurrentMap<String, Connector> tempConnectors = new ConcurrentHashMap<>();
private final ConcurrentMap<String, IsolatedConnector<?>> tempConnectors = new ConcurrentHashMap<>();

public AbstractHerder(Worker worker,
String workerId,
Expand Down Expand Up @@ -389,13 +393,13 @@ public ConnectorStateInfo.TaskState taskStatus(ConnectorTaskId id) {
status.workerId(), status.trace());
}

protected Map<String, ConfigValue> validateSinkConnectorConfig(SinkConnector connector, ConfigDef configDef, Map<String, String> config) {
protected Map<String, ConfigValue> validateSinkConnectorConfig(IsolatedSinkConnector connector, ConfigDef configDef, Map<String, String> config) {
Map<String, ConfigValue> result = configDef.validateAll(config);
SinkConnectorConfig.validate(config, result);
return result;
}

protected Map<String, ConfigValue> validateSourceConnectorConfig(SourceConnector connector, ConfigDef configDef, Map<String, String> config) {
protected Map<String, ConfigValue> validateSourceConnectorConfig(IsolatedSourceConnector connector, ConfigDef configDef, Map<String, String> config) {
return configDef.validateAll(config);
}

Expand Down Expand Up @@ -647,7 +651,7 @@ ConfigInfos validateConnectorConfig(
Map<String, String> connectorProps,
Function<String, TemporaryStage> reportStage,
boolean doLog
) {
) throws Exception {
String stageDescription;
if (worker.configTransformer() != null) {
stageDescription = "resolving transformed configuration properties for the connector";
Expand All @@ -659,25 +663,26 @@ ConfigInfos validateConnectorConfig(
if (connType == null)
throw new BadRequestException("Connector config " + connectorProps + " contains no connector type");

Connector connector = getConnector(connType);
IsolatedConnector<?> connector = getConnector(connType);
ClassLoader connectorLoader = plugins().connectorLoader(connType);
try (LoaderSwap loaderSwap = plugins().withClassLoader(connectorLoader)) {
org.apache.kafka.connect.health.ConnectorType connectorType;
ConfigDef enrichedConfigDef;
Map<String, ConfigValue> validatedConnectorConfig;
if (connector instanceof SourceConnector) {
PluginType type = connector.type();
if (type == PluginType.SOURCE) {
connectorType = org.apache.kafka.connect.health.ConnectorType.SOURCE;
enrichedConfigDef = ConnectorConfig.enrich(plugins(), SourceConnectorConfig.configDef(), connectorProps, false);
stageDescription = "validating source connector-specific properties for the connector";
try (TemporaryStage stage = reportStage.apply(stageDescription)) {
validatedConnectorConfig = validateSourceConnectorConfig((SourceConnector) connector, enrichedConfigDef, connectorProps);
validatedConnectorConfig = validateSourceConnectorConfig((IsolatedSourceConnector) connector, enrichedConfigDef, connectorProps);
}
} else {
connectorType = org.apache.kafka.connect.health.ConnectorType.SINK;
enrichedConfigDef = ConnectorConfig.enrich(plugins(), SinkConnectorConfig.configDef(), connectorProps, false);
stageDescription = "validating sink connector-specific properties for the connector";
try (TemporaryStage stage = reportStage.apply(stageDescription)) {
validatedConnectorConfig = validateSinkConnectorConfig((SinkConnector) connector, enrichedConfigDef, connectorProps);
validatedConnectorConfig = validateSinkConnectorConfig((IsolatedSinkConnector) connector, enrichedConfigDef, connectorProps);
}
}

Expand All @@ -703,7 +708,7 @@ ConfigInfos validateConnectorConfig(
throw new BadRequestException(
String.format(
"%s.config() must return a ConfigDef that is not null.",
connector.getClass().getName()
connector.pluginClass().getName()
)
);
}
Expand All @@ -717,7 +722,7 @@ ConfigInfos validateConnectorConfig(
throw new BadRequestException(
String.format(
"%s.validate() must return a Config that is not null.",
connector.getClass().getName()
connector.pluginClass().getName()
)
);
}
Expand Down Expand Up @@ -757,7 +762,7 @@ ConfigInfos validateConnectorConfig(
ConnectorConfig.CONNECTOR_CLIENT_PRODUCER_OVERRIDES_PREFIX,
connectorConfig,
ProducerConfig.configDef(),
connector.getClass(),
connector.pluginClass(),
connectorType,
ConnectorClientConfigRequest.ClientType.PRODUCER,
connectorClientConfigOverridePolicy);
Expand All @@ -771,7 +776,7 @@ ConfigInfos validateConnectorConfig(
ConnectorConfig.CONNECTOR_CLIENT_ADMIN_OVERRIDES_PREFIX,
connectorConfig,
AdminClientConfig.configDef(),
connector.getClass(),
connector.pluginClass(),
connectorType,
ConnectorClientConfigRequest.ClientType.ADMIN,
connectorClientConfigOverridePolicy);
Expand All @@ -785,7 +790,7 @@ ConfigInfos validateConnectorConfig(
ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX,
connectorConfig,
ConsumerConfig.configDef(),
connector.getClass(),
connector.pluginClass(),
connectorType,
ConnectorClientConfigRequest.ClientType.CONSUMER,
connectorClientConfigOverridePolicy);
Expand Down Expand Up @@ -945,7 +950,7 @@ private static ConfigValueInfo convertConfigValue(ConfigValue configValue, Type
return new ConfigValueInfo(configValue.name(), value, recommendedValues, configValue.errorMessages(), configValue.visible());
}

protected Connector getConnector(String connType) {
protected IsolatedConnector<?> getConnector(String connType) {
return tempConnectors.computeIfAbsent(connType, k -> plugins().newConnector(k));
}

Expand All @@ -964,7 +969,7 @@ public ConnectorType connectorType(Map<String, String> connConfig) {
return ConnectorType.UNKNOWN;
}
try {
return ConnectorType.from(getConnector(connClass).getClass());
return ConnectorType.from(getConnector(connClass).pluginClass());
} catch (ConnectException e) {
log.warn("Unable to retrieve connector type", e);
return ConnectorType.UNKNOWN;
Expand Down
Loading