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);
+        }
+    }
+}

Reply via email to