Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
1 change: 1 addition & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -1274,6 +1274,7 @@ project(':clients') {

testImplementation libs.bcpkix
testImplementation libs.junitJupiter
testImplementation libs.log4j
Comment thread
C0urante marked this conversation as resolved.
Outdated
testImplementation libs.mockitoInline

testRuntimeOnly libs.slf4jlog4j
Expand Down
5 changes: 1 addition & 4 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,7 @@

<subpackage name="utils">
<allow pkg="org.apache.kafka.common" />
<allow pkg="org.apache.log4j" />
</subpackage>

<subpackage name="quotas">
Expand Down Expand Up @@ -454,9 +455,6 @@
<allow pkg="com.fasterxml.jackson" />
<allow pkg="kafka.utils" />
<allow pkg="org.apache.zookeeper" />
<subpackage name="testutil">
<allow pkg="org.apache.log4j" />
</subpackage>
</subpackage>
</subpackage>
</subpackage>
Expand Down Expand Up @@ -573,7 +571,6 @@
<allow pkg="com.fasterxml.jackson" />
<allow pkg="org.apache.http"/>
<allow pkg="io.swagger.v3.oas.annotations"/>
<allow pkg="kafka.utils" />
<subpackage name="resources">
<allow pkg="org.apache.log4j" />
</subpackage>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.kafka.streams.processor.internals.testutil;
package org.apache.kafka.common.utils;

import org.apache.log4j.AppenderSkeleton;
import org.apache.log4j.Level;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,13 @@ public void remove(ConnectorTaskId id) {
}
}

private void commit(WorkerSourceTask workerTask) {
// Visible for testing
static void commit(WorkerSourceTask workerTask) {
if (!workerTask.shouldCommitOffsets()) {
Comment thread
C0urante marked this conversation as resolved.
Outdated
log.trace("{} Skipping offset commit as there are no offsets that should be committed", workerTask);
return;
}

log.debug("{} Committing offsets", workerTask);
try {
if (workerTask.commitOffsets()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,14 @@ protected void finalOffsetCommit(boolean failed) {
commitOffsets();
}

/**
* @return whether an attempt to commit offsets should be made for the task (i.e., there are pending uncommitted
* offsets and the task's producer has not already failed to send a record with a non-retriable error).
*/
public boolean shouldCommitOffsets() {
Comment thread
tombentley marked this conversation as resolved.
Outdated
return !isFailed();
}

public boolean commitOffsets() {
long commitTimeoutMs = workerConfig.getLong(WorkerConfig.OFFSET_COMMIT_TIMEOUT_MS_CONFIG);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ abstract class WorkerTask implements Runnable {
private final CountDownLatch shutdownLatch = new CountDownLatch(1);
private final TaskMetricsGroup taskMetricsGroup;
private volatile TargetState targetState;
private volatile boolean failed;
private volatile boolean stopping; // indicates whether the Worker has asked the task to stop
private volatile boolean cancelled; // indicates whether the Worker has cancelled the task (e.g. because of slow shutdown)
private final ErrorHandlingMetrics errorMetrics;
Expand All @@ -84,6 +85,7 @@ public WorkerTask(ConnectorTaskId id,
this.statusListener = taskMetricsGroup;
this.loader = loader;
this.targetState = initialState;
this.failed = false;
this.stopping = false;
this.cancelled = false;
this.taskMetricsGroup.recordState(this.targetState);
Expand Down Expand Up @@ -163,6 +165,10 @@ public void removeMetrics() {

protected abstract void close();

protected boolean isFailed() {
return failed;
}

protected boolean isStopping() {
return stopping;
}
Expand Down Expand Up @@ -196,6 +202,7 @@ private void doRun() throws InterruptedException {
statusListener.onStartup(id);
execute();
} catch (Throwable t) {
failed = true;
if (cancelled) {
log.warn("{} After being scheduled for shutdown, the orphan task threw an uncaught exception. A newer instance of this task might be already running", this, t);
} else if (stopping) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import org.apache.kafka.connect.util.ParameterizedTest;
import org.apache.kafka.connect.util.TopicAdmin;
import org.apache.kafka.connect.util.TopicCreationGroup;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.easymock.Capture;
import org.easymock.EasyMock;
import org.easymock.IAnswer;
Expand Down Expand Up @@ -99,6 +100,7 @@
import static org.apache.kafka.connect.runtime.TopicCreationConfig.REPLICATION_FACTOR_CONFIG;
import static org.apache.kafka.connect.runtime.WorkerConfig.TOPIC_CREATION_ENABLE_CONFIG;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertThrows;
import static org.junit.Assert.assertTrue;
Expand Down Expand Up @@ -372,7 +374,7 @@ public void testFailureInPoll() throws Exception {

sourceTask.stop();
EasyMock.expectLastCall();
expectOffsetFlush(true);
expectEmptyOffsetFlush();

expectClose();

Expand All @@ -382,8 +384,9 @@ public void testFailureInPoll() throws Exception {
Future<?> taskFuture = executor.submit(workerTask);

assertTrue(awaitLatch(pollLatch));
//Failure in poll should trigger automatic stop of the worker
//Failure in poll should trigger automatic stop of the task
assertTrue(workerTask.awaitStop(1000));
assertShouldSkipCommit();

taskFuture.get();
assertPollMetrics(0);
Expand Down Expand Up @@ -467,6 +470,7 @@ public void testFailureInPollAfterStop() throws Exception {
workerTask.stop();
workerStopLatch.countDown();
assertTrue(workerTask.awaitStop(1000));
assertShouldSkipCommit();

taskFuture.get();
assertPollMetrics(0);
Expand All @@ -483,11 +487,11 @@ public void testPollReturnsNoRecords() throws Exception {

// We'll wait for some data, then trigger a flush
final CountDownLatch pollLatch = expectEmptyPolls(1, new AtomicInteger());
expectOffsetFlush(true);
expectEmptyOffsetFlush();

sourceTask.stop();
EasyMock.expectLastCall();
expectOffsetFlush(true);
expectEmptyOffsetFlush();

statusListener.onShutdown(taskId);
EasyMock.expectLastCall();
Expand Down Expand Up @@ -528,7 +532,7 @@ public void testCommit() throws Exception {

sourceTask.stop();
EasyMock.expectLastCall();
expectOffsetFlush(true);
expectEmptyOffsetFlush();

statusListener.onShutdown(taskId);
EasyMock.expectLastCall();
Expand Down Expand Up @@ -647,10 +651,6 @@ public void testSendRecordsProducerCallbackFail() throws Exception {

@Test
public void testSendRecordsProducerSendFailsImmediately() {
if (!enableTopicCreation)
// should only test with topic creation enabled
return;

createWorkerTask();

SourceRecord record1 = new SourceRecord(PARTITION, OFFSET, TOPIC, 1, KEY_SCHEMA, KEY, RECORD_SCHEMA, RECORD);
Expand Down Expand Up @@ -1023,6 +1023,12 @@ private void expectOffsetFlush(boolean succeed) throws Exception {
}
}

private void expectEmptyOffsetFlush() throws Exception {
EasyMock.expect(offsetWriter.beginFlush()).andReturn(false);
sourceTask.commit();
EasyMock.expectLastCall();
}

private void assertPollMetrics(int minimumPollCountExpected) {
MetricGroup sourceTaskGroup = workerTask.sourceTaskMetricsGroup().metricGroup();
MetricGroup taskGroup = workerTask.taskMetricsGroup().metricGroup();
Expand Down Expand Up @@ -1109,4 +1115,20 @@ private void expectTopicCreation(String topic) {
EasyMock.expect(admin.createOrFindTopics(EasyMock.capture(newTopicCapture))).andReturn(createdTopic(topic));
}
}

private void assertShouldSkipCommit() {
assertFalse(workerTask.shouldCommitOffsets());

LogCaptureAppender.setClassLoggerToTrace(SourceTaskOffsetCommitter.class);
LogCaptureAppender.setClassLoggerToTrace(WorkerSourceTask.class);
try (LogCaptureAppender committerAppender = LogCaptureAppender.createAndRegister(SourceTaskOffsetCommitter.class);
LogCaptureAppender taskAppender = LogCaptureAppender.createAndRegister(WorkerSourceTask.class)) {
SourceTaskOffsetCommitter.commit(workerTask);
assertEquals(Collections.emptyList(), taskAppender.getMessages());

List<String> committerMessages = committerAppender.getMessages();
assertEquals(1, committerMessages.size());
assertTrue(committerMessages.get(0).contains("Skipping offset commit"));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@
*/
package org.apache.kafka.connect.runtime.rest;

import kafka.utils.LogCaptureAppender;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.http.HttpHost;
Expand All @@ -31,6 +30,7 @@
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.connect.rest.ConnectRestExtension;
import org.apache.kafka.connect.runtime.Herder;
import org.apache.kafka.connect.runtime.WorkerConfig;
Expand Down Expand Up @@ -362,10 +362,10 @@ public void testRequestLogs() throws IOException, InterruptedException {
// Stop the server to flush all logs
server.stop();

Collection<String> logMessages = restServerAppender.getRenderedMessages();
Collection<String> logMessages = restServerAppender.getMessages();
LogCaptureAppender.unregister(restServerAppender);
restServerAppender.close();
String expectedlogContent = "\"GET / HTTP/1.1\" " + String.valueOf(response.getStatusLine().getStatusCode());
String expectedlogContent = "\"GET / HTTP/1.1\" " + response.getStatusLine().getStatusCode();
assertTrue(logMessages.stream().anyMatch(logMessage -> logMessage.contains(expectedlogContent)));
}

Expand Down
6 changes: 0 additions & 6 deletions core/src/test/scala/unit/kafka/utils/LogCaptureAppender.scala
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@ package kafka.utils
import org.apache.log4j.{AppenderSkeleton, Level, Logger}
import org.apache.log4j.spi.LoggingEvent

import scala.jdk.CollectionConverters._
import scala.collection.mutable.ListBuffer

class LogCaptureAppender extends AppenderSkeleton {
Expand All @@ -38,11 +37,6 @@ class LogCaptureAppender extends AppenderSkeleton {
}
}

def getRenderedMessages: java.util.List[String] = {
return getMessages.map(e => e.getRenderedMessage).asJava
}


override def close(): Unit = {
events.synchronized {
events.clear()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@
import org.apache.kafka.streams.processor.internals.ThreadMetadataImpl;
import org.apache.kafka.streams.processor.internals.TopologyMetadata;
import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.state.Stores;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@
import org.apache.kafka.streams.processor.FailOnInvalidTimestamp;
import org.apache.kafka.streams.processor.TimestampExtractor;
import org.apache.kafka.streams.processor.internals.StreamsPartitionAssignor;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.junit.Before;
import org.junit.Rule;
import org.junit.Test;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.errors.TimeoutException;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
Expand All @@ -31,7 +32,6 @@
import org.apache.kafka.streams.kstream.Transformer;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.processor.PunctuationType;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.test.TestUtils;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
import org.apache.kafka.streams.kstream.TimeWindows;
import org.apache.kafka.streams.kstream.Windowed;
import org.apache.kafka.streams.kstream.Windows;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.state.SessionStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@
import org.apache.kafka.streams.kstream.StreamJoined;
import org.apache.kafka.streams.processor.internals.InternalTopicConfig;
import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.TestInputTopic;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.state.WindowBytesStoreSupplier;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.KeyValueTimestamp;
import org.apache.kafka.streams.StreamsBuilder;
Expand All @@ -43,7 +44,6 @@
import org.apache.kafka.streams.kstream.Joined;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.test.MockApiProcessor;
import org.apache.kafka.test.MockApiProcessorSupplier;
import org.apache.kafka.test.MockValueJoiner;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,8 @@
import org.apache.kafka.streams.processor.internals.metrics.TaskMetrics;
import org.apache.kafka.streams.processor.internals.ProcessorRecordContext;
import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender.Event;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender.Event;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.apache.kafka.streams.state.SessionBytesStoreSupplier;
import org.apache.kafka.streams.state.SessionStore;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,8 @@
import org.apache.kafka.streams.kstream.SlidingWindows;
import org.apache.kafka.streams.kstream.TimeWindowedDeserializer;
import org.apache.kafka.streams.kstream.Windowed;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender.Event;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender.Event;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.state.ValueAndTimestamp;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.KeyValueTimestamp;
Expand All @@ -47,7 +48,6 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.processor.internals.ProcessorNode;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.state.TimestampedWindowStore;
import org.apache.kafka.streams.state.WindowBytesStoreSupplier;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@
import org.apache.kafka.streams.processor.api.MockProcessorContext;
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.test.TestRecord;
import org.apache.kafka.test.MockApiProcessor;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
import org.apache.kafka.streams.processor.api.MockProcessorContext;
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender;
import org.apache.kafka.common.utils.LogCaptureAppender;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.test.TestRecord;
import org.apache.kafka.test.MockApiProcessor;
Expand Down
Loading