Copilot commented on code in PR #23422:
URL: https://github.com/apache/kafka/pull/23422#discussion_r3979411620
##########
connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java:
##########
@@ -1920,6 +1921,10 @@ public WorkerTask<ConsumerRecord<byte[], byte[]>,
SinkRecord> doBuild(
Map<String, Object> consumerProps = baseConsumerConfigs(
id.connector(), "connector-consumer-" + id, config,
connectorConfig, connectorClass,
connectorClientConfigOverridePolicy, kafkaClusterId,
ConnectorType.SINK);
+ Object groupProtocol =
consumerProps.get(ConsumerConfig.GROUP_PROTOCOL_CONFIG);
+ if (groupProtocol != null &&
GroupProtocol.CONSUMER.name().equalsIgnoreCase(groupProtocol.toString())) {
+ log.warn("Sink task {} uses group.protocol=CONSUMER, which has
not been fully tested for production use with Kafka Connect.", id);
+ }
Review Comment:
This WARN will be emitted on every sink task start, which can become noisy
in environments with frequent task restarts/rebalances. Consider logging once
per worker/connector (e.g., track warned connector/task IDs in a concurrent
set, or a single AtomicBoolean if worker-level is sufficient), or use a
rate-limited logger if the project has one.
##########
connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java:
##########
@@ -711,6 +712,80 @@ public void testAddRemoveSinkTask(boolean
enableTopicCreation) {
verifyExecutorSubmit();
}
+ @ParameterizedTest
+ @CsvSource({
+ "true, consumer, true",
+ "false, consumer, true",
+ "true, classic, false",
+ "false, classic, false"
+ })
+ public void testWarnWhenStartingSinkTaskWithConsumerGroupProtocol(
+ boolean enableTopicCreation,
+ String groupProtocol,
+ boolean expectWarning) {
+ setup(enableTopicCreation);
+ workerProps.put("consumer." + ConsumerConfig.GROUP_PROTOCOL_CONFIG,
groupProtocol);
+ config = new StandaloneConfig(workerProps);
+ // Most of the other cases use source tasks; we make sure to get code
coverage for sink tasks here as well
+ SinkTask task = mock(TestSinkTask.class);
+ mockKafkaClusterId();
+ mockVersionedTaskIsolation(SampleSinkConnector.class,
TestSinkTask.class, null, sinkConnector, task);
+
mockVersionedTaskConverterFromConnector(ConnectorConfig.KEY_CONVERTER_CLASS_CONFIG,
ConnectorConfig.KEY_CONVERTER_VERSION_CONFIG, taskKeyConverter);
+
mockVersionedTaskConverterFromConnector(ConnectorConfig.VALUE_CONVERTER_CLASS_CONFIG,
ConnectorConfig.VALUE_CONVERTER_VERSION_CONFIG, taskValueConverter);
+ mockVersionedTaskHeaderConverterFromConnector(taskHeaderConverter);
+ mockExecutorFakeSubmit(WorkerTask.class);
+
+ Map<String, String> origProps = Map.of(TaskConfig.TASK_CLASS_CONFIG,
TestSinkTask.class.getName());
+
+ worker = new Worker(WORKER_ID, new MockTime(), plugins, config,
offsetBackingStore, executorService,
+ noneConnectorClientConfigOverridePolicy, null);
+ worker.herder = herder;
+ worker.start();
+
+ assertStatistics(worker, 0, 0);
+ assertEquals(Set.of(), worker.taskIds());
+ Map<String, String> connectorConfigs = anyConnectorConfigMap();
+ connectorConfigs.put(TOPICS_CONFIG, "t1");
+ connectorConfigs.put(CONNECTOR_CLASS_CONFIG,
SampleSinkConnector.class.getName());
+
+ ClusterConfigState configState = new ClusterConfigState(
+ 0,
+ null,
+ Map.of(CONNECTOR_ID, 1),
+ Map.of(CONNECTOR_ID, connectorConfigs),
+ Map.of(CONNECTOR_ID, TargetState.STARTED),
+ Map.of(TASK_ID, origProps),
+ Map.of(),
+ Map.of(),
+ Map.of(CONNECTOR_ID, new
AppliedConnectorConfig(connectorConfigs)),
+ Set.of(),
+ Set.of()
+ );
+ boolean warningLogged;
+ try (LogCaptureAppender appender =
LogCaptureAppender.createAndRegister(Worker.class)) {
+ assertTrue(worker.startSinkTask(TASK_ID, configState,
connectorConfigs, origProps, taskStatusListener, TargetState.STARTED));
+ warningLogged =
appender.getMessages("WARN").stream().anyMatch(message ->
message.contains("group.protocol=CONSUMER")
+ && message.contains("not been fully tested for production")
+ && message.contains(TASK_ID.toString()));
Review Comment:
This assertion is tightly coupled to the exact log phrasing and case (e.g.,
the literal `group.protocol=CONSUMER` and substring text). To reduce
brittleness, consider asserting on a more stable signature (e.g., only
`TASK_ID` + a shorter invariant phrase), or centralize the warning message as a
constant in `Worker` so the test can reference it indirectly.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]