This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11986-fa66cce41587838581014dca039e8c249adeb219 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit be265d45c6fd586a17dc97f8427219ed1d8b302d Author: zhiwei.niu <[email protected]> AuthorDate: Sun Aug 30 04:23:29 2026 +0000 [Feature][Zeta] Add intermediate queue comparison benchmark (#11986) --- .github/workflows/benchmarks.yml | 1 + seatunnel-benchmarks/pom.xml | 4 + .../benchmark/IntermediateQueueBenchmark.java | 120 +++++++++ .../benchmark/IntermediateQueueBenchmarkState.java | 274 +++++++++++++++++++++ .../benchmark/IntermediateQueueBenchmarkTest.java | 60 +++++ 5 files changed, 459 insertions(+) diff --git a/.github/workflows/benchmarks.yml b/.github/workflows/benchmarks.yml index 25f74cb37c..c18bc9fe6d 100644 --- a/.github/workflows/benchmarks.yml +++ b/.github/workflows/benchmarks.yml @@ -35,6 +35,7 @@ on: options: - '.*' - 'SeaTunnelRowBenchmark' + - 'IntermediateQueueBenchmark' - 'SeaTunnelPipelineBenchmark' - 'sourceSink$' - 'sourceTransformSink$' diff --git a/seatunnel-benchmarks/pom.xml b/seatunnel-benchmarks/pom.xml index a9788e1dba..b0544df0ff 100644 --- a/seatunnel-benchmarks/pom.xml +++ b/seatunnel-benchmarks/pom.xml @@ -60,6 +60,10 @@ </exclusion> </exclusions> </dependency> + <dependency> + <groupId>com.lmax</groupId> + <artifactId>disruptor</artifactId> + </dependency> <dependency> <groupId>org.apache.seatunnel</groupId> <artifactId>seatunnel-shade-hazelcast</artifactId> diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmark.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmark.java new file mode 100644 index 0000000000..7ce77949db --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmark.java @@ -0,0 +1,120 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.benchmark; + +import org.apache.seatunnel.engine.common.config.server.QueueType; + +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.Threads; +import org.openjdk.jmh.runner.Runner; +import org.openjdk.jmh.runner.RunnerException; +import org.openjdk.jmh.runner.options.Options; +import org.openjdk.jmh.runner.options.OptionsBuilder; +import org.openjdk.jmh.runner.options.VerboseMode; + +import java.util.concurrent.TimeUnit; + +/** Compares record handoff throughput of the two production intermediate queue implementations. */ +@OutputTimeUnit(TimeUnit.SECONDS) +@Threads(1) +public class IntermediateQueueBenchmark extends BenchmarkBase { + + public static void main(String[] args) throws RunnerException { + Options options = + new OptionsBuilder() + .verbosity(VerboseMode.NORMAL) + .include(".*" + IntermediateQueueBenchmark.class.getCanonicalName() + ".*") + .build(); + new Runner(options).run(); + } + + @Benchmark + public long blockingQueueRecordHandoff(BlockingQueueState state) { + return state.publish(); + } + + @Benchmark + public long disruptorRecordHandoff(DisruptorQueueState state) { + return state.publish(); + } + + @State(Scope.Thread) + public static class BlockingQueueState { + + @Param({"1024"}) + private int capacity; + + @Param({"4096"}) + private int recordPoolSize; + + private IntermediateQueueBenchmarkState delegate; + + @Setup(Level.Trial) + public void setUp() throws Exception { + delegate = + new IntermediateQueueBenchmarkState( + QueueType.BLOCKINGQUEUE, capacity, recordPoolSize); + delegate.setUp(); + } + + long publish() { + return delegate.publish(); + } + + @TearDown(Level.Trial) + public void tearDown() throws Exception { + delegate.tearDown(); + } + } + + @State(Scope.Thread) + public static class DisruptorQueueState { + + @Param({"1024"}) + private int capacity; + + @Param({"4096"}) + private int recordPoolSize; + + private IntermediateQueueBenchmarkState delegate; + + @Setup(Level.Trial) + public void setUp() throws Exception { + delegate = + new IntermediateQueueBenchmarkState( + QueueType.DISRUPTOR, capacity, recordPoolSize); + delegate.setUp(); + } + + long publish() { + return delegate.publish(); + } + + @TearDown(Level.Trial) + public void tearDown() throws Exception { + delegate.tearDown(); + } + } +} diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkState.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkState.java new file mode 100644 index 0000000000..24720fc953 --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkState.java @@ -0,0 +1,274 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.benchmark; + +import org.apache.seatunnel.api.common.metrics.Counter; +import org.apache.seatunnel.api.common.metrics.ThreadSafeCounter; +import org.apache.seatunnel.api.table.type.Record; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +import org.apache.seatunnel.api.transform.Collector; +import org.apache.seatunnel.engine.common.config.server.QueueType; +import org.apache.seatunnel.engine.common.utils.concurrent.CompletableFuture; +import org.apache.seatunnel.engine.server.task.flow.IntermediateQueueFlowLifeCycle; +import org.apache.seatunnel.engine.server.task.group.queue.AbstractIntermediateQueue; +import org.apache.seatunnel.engine.server.task.group.queue.IntermediateBlockingQueue; +import org.apache.seatunnel.engine.server.task.group.queue.IntermediateDisruptor; +import org.apache.seatunnel.engine.server.task.group.queue.disruptor.RecordEvent; +import org.apache.seatunnel.engine.server.task.group.queue.disruptor.RecordEventFactory; + +import com.lmax.disruptor.YieldingWaitStrategy; +import com.lmax.disruptor.dsl.Disruptor; +import com.lmax.disruptor.dsl.ProducerType; +import com.lmax.disruptor.util.DaemonThreadFactory; + +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.LockSupport; + +/** Owns the queue lifecycle and reusable records shared by the queue comparison benchmarks. */ +final class IntermediateQueueBenchmarkState { + + private static final int PAYLOAD_SIZE = 256; + private static final long CONSUMER_DRAIN_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(10); + + private final QueueType queueType; + private final int capacity; + private final int recordPoolSize; + private final AtomicLong consumedRecords = new AtomicLong(); + private final AtomicLong consumedChecksum = new AtomicLong(); + private final AtomicReference<Throwable> consumerFailure = new AtomicReference<>(); + + private Record<?>[] records; + private AbstractIntermediateQueue<?> queue; + private IntermediateQueueFlowLifeCycle<?> flowLifeCycle; + private Thread blockingQueueConsumer; + private long publishedRecords; + + IntermediateQueueBenchmarkState(QueueType queueType, int capacity, int recordPoolSize) { + this.queueType = queueType; + this.capacity = capacity; + this.recordPoolSize = recordPoolSize; + } + + void setUp() throws Exception { + validateParameters(); + records = createRecords(recordPoolSize); + queue = createQueue(); + flowLifeCycle = + new IntermediateQueueFlowLifeCycle<>(null, new CompletableFuture<>(), queue); + + if (queueType == QueueType.BLOCKINGQUEUE) { + startBlockingQueueConsumer(); + } else { + flowLifeCycle.collect(new BenchmarkCollector()); + } + } + + long publish() { + checkConsumerFailure(); + Record<?> record = records[(int) (publishedRecords & (recordPoolSize - 1))]; + flowLifeCycle.received(record); + return ++publishedRecords; + } + + void tearDown() throws Exception { + Throwable failure = null; + try { + awaitConsumerDrain(); + checkConsumerFailure(); + } catch (Throwable throwable) { + failure = throwable; + } + + if (blockingQueueConsumer != null) { + blockingQueueConsumer.interrupt(); + try { + blockingQueueConsumer.join(TimeUnit.SECONDS.toMillis(1)); + if (blockingQueueConsumer.isAlive()) { + throw new IllegalStateException("Blocking queue consumer did not stop"); + } + } catch (Throwable throwable) { + failure = combineFailures(failure, throwable); + } + } + + if (flowLifeCycle != null) { + try { + flowLifeCycle.close(); + } catch (Throwable throwable) { + failure = combineFailures(failure, throwable); + } + } + + if (failure != null) { + rethrow(failure); + } + } + + long getPublishedRecords() { + return publishedRecords; + } + + long getConsumedRecords() { + return consumedRecords.get(); + } + + long getConsumedChecksum() { + return consumedChecksum.get(); + } + + private void validateParameters() { + if (capacity <= 0 || Integer.bitCount(capacity) != 1) { + throw new IllegalArgumentException("capacity must be a positive power of two"); + } + // Keep at least one reusable record per queue slot so a saturated queue does not contain + // duplicate record instances. Records remain read-only after trial setup. + if (recordPoolSize < capacity || Integer.bitCount(recordPoolSize) != 1) { + throw new IllegalArgumentException( + "recordPoolSize must be a power of two and at least capacity"); + } + } + + private AbstractIntermediateQueue<?> createQueue() { + Counter totalQueueSize = new ThreadSafeCounter("totalQueueSize"); + Counter queueSize = new ThreadSafeCounter("queueSize"); + Counter putBlockedNs = new ThreadSafeCounter("putBlockedNs"); + Counter flushSuccess = new ThreadSafeCounter("flushSuccess"); + Counter flushFailure = new ThreadSafeCounter("flushFailure"); + + if (queueType == QueueType.BLOCKINGQUEUE) { + return new IntermediateBlockingQueue( + new ArrayBlockingQueue<>(capacity), + totalQueueSize, + queueSize, + putBlockedNs, + flushSuccess, + flushFailure); + } + + Disruptor<RecordEvent> disruptor = + new Disruptor<>( + new RecordEventFactory(), + capacity, + DaemonThreadFactory.INSTANCE, + ProducerType.SINGLE, + new YieldingWaitStrategy()); + return new IntermediateDisruptor( + disruptor, totalQueueSize, queueSize, putBlockedNs, flushSuccess, flushFailure); + } + + private static Record<?>[] createRecords(int size) { + Record<?>[] recordPool = new Record<?>[size]; + for (int i = 0; i < size; i++) { + byte[] payload = new byte[PAYLOAD_SIZE]; + for (int j = 0; j < payload.length; j++) { + payload[j] = (byte) (i + j); + } + recordPool[i] = + new Record<>( + new SeaTunnelRow( + new Object[] { + (long) i, "queue-benchmark-" + i, payload, i % 128 + })); + } + return recordPool; + } + + private void startBlockingQueueConsumer() { + blockingQueueConsumer = + new Thread( + () -> { + try { + while (!Thread.currentThread().isInterrupted()) { + flowLifeCycle.collect(new BenchmarkCollector()); + } + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + } catch (Throwable throwable) { + consumerFailure.compareAndSet(null, throwable); + } + }, + "intermediate-blocking-queue-benchmark-consumer"); + blockingQueueConsumer.setDaemon(true); + blockingQueueConsumer.start(); + } + + private void consume(Record<?> record) { + SeaTunnelRow row = (SeaTunnelRow) record.getData(); + byte[] payload = (byte[]) row.getField(2); + long checksum = + (Long) row.getField(0) + + ((String) row.getField(1)).length() + + Byte.toUnsignedInt(payload[0]) + + Byte.toUnsignedInt(payload[payload.length - 1]) + + (Integer) row.getField(3); + consumedChecksum.addAndGet(checksum); + consumedRecords.incrementAndGet(); + } + + private void awaitConsumerDrain() { + long deadline = System.nanoTime() + CONSUMER_DRAIN_TIMEOUT_NANOS; + while (consumedRecords.get() != publishedRecords && System.nanoTime() < deadline) { + checkConsumerFailure(); + LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(1)); + } + if (consumedRecords.get() != publishedRecords) { + throw new IllegalStateException( + String.format( + "Timed out waiting for %s: published=%d, consumed=%d", + queueType, publishedRecords, consumedRecords.get())); + } + } + + private void checkConsumerFailure() { + Throwable failure = consumerFailure.get(); + if (failure != null) { + throw new IllegalStateException("Queue consumer failed", failure); + } + } + + private static Throwable combineFailures(Throwable first, Throwable second) { + if (first == null) { + return second; + } + first.addSuppressed(second); + return first; + } + + private static void rethrow(Throwable failure) throws Exception { + if (failure instanceof Exception) { + throw (Exception) failure; + } + throw (Error) failure; + } + + private final class BenchmarkCollector implements Collector<Record<?>> { + + @Override + public void collect(Record<?> record) { + consume(record); + } + + @Override + public void close() { + // Nothing to close. + } + } +} diff --git a/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkTest.java b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkTest.java new file mode 100644 index 0000000000..87cbcea545 --- /dev/null +++ b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/IntermediateQueueBenchmarkTest.java @@ -0,0 +1,60 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.benchmark; + +import org.apache.seatunnel.engine.common.config.server.QueueType; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class IntermediateQueueBenchmarkTest { + + private static final int RECORD_COUNT = 10_000; + + @Test + void shouldTransferAllRecordsThroughEachQueueImplementation() throws Exception { + for (QueueType queueType : QueueType.values()) { + IntermediateQueueBenchmarkState state = + new IntermediateQueueBenchmarkState(queueType, 1024, 4096); + state.setUp(); + try { + for (int i = 0; i < RECORD_COUNT; i++) { + state.publish(); + } + } finally { + state.tearDown(); + } + + assertEquals(RECORD_COUNT, state.getPublishedRecords()); + assertEquals(RECORD_COUNT, state.getConsumedRecords()); + assertEquals(21_783_822L, state.getConsumedChecksum()); + } + } + + @Test + void shouldRejectNonPowerOfTwoCapacityForEachQueueType() { + for (QueueType queueType : QueueType.values()) { + IntermediateQueueBenchmarkState state = + new IntermediateQueueBenchmarkState(queueType, 1000, 4096); + + assertThrows(IllegalArgumentException.class, state::setUp); + } + } +}
