From 985386989d88894da7758ef9cef4ee10009933f0 Mon Sep 17 00:00:00 2001 From: wwang Date: Mon, 13 Jul 2026 11:18:44 +0800 Subject: [PATCH 01/13] fix(QTDI-2709): Add ProcessorBufferingTest to evaluate memory consumption during record emission --- .../components/ProcessorBufferingTest.java | 372 ++++++++++++++++++ 1 file changed, 372 insertions(+) create mode 100644 component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java new file mode 100644 index 0000000000000..299b97edfc068 --- /dev/null +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java @@ -0,0 +1,372 @@ +/** + * Copyright (C) 2006-2026 Talend Inc. - www.talend.com + * + * Licensed 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.talend.sdk.component.runtime.di.beam.components; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.File; +import java.io.ObjectInputStream; +import java.io.ObjectOutputStream; +import java.io.Serializable; +import java.util.HashMap; +import java.util.Map; +import java.util.stream.Stream; + +import javax.json.bind.Jsonb; + +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.talend.sdk.component.api.processor.AfterGroup; +import org.talend.sdk.component.api.processor.BeforeGroup; +import org.talend.sdk.component.api.processor.ElementListener; +import org.talend.sdk.component.api.processor.Input; +import org.talend.sdk.component.api.processor.LastGroup; +import org.talend.sdk.component.api.processor.Output; +import org.talend.sdk.component.api.processor.OutputEmitter; +import org.talend.sdk.component.api.processor.Processor; +import org.talend.sdk.component.api.record.Record; +import org.talend.sdk.component.api.service.record.RecordBuilderFactory; +import org.talend.sdk.component.runtime.di.AutoChunkProcessor; +import org.talend.sdk.component.runtime.di.InputsHandler; +import org.talend.sdk.component.runtime.di.JobStateAware; +import org.talend.sdk.component.runtime.di.OutputsHandler; +import org.talend.sdk.component.runtime.manager.ComponentManager; +import org.talend.sdk.component.runtime.output.InputFactory; +import org.talend.sdk.component.runtime.output.OutputFactory; +import org.talend.sdk.component.runtime.record.RecordConverters; + +import lombok.Getter; +import lombok.ToString; + +/** + * Manual test to observe processor output buffering problem (QTDI-2709). + * + * Run this test to see memory consumption when a processor emits many records. + * With the current push model, ALL records are buffered in memory before the + * consumer can drain them. + * + * Expected behavior on current code: test FAILS (memory exceeds threshold). + * After iterator pattern is implemented: test PASSES (memory stays low). + */ +class ProcessorBufferingTest { + + protected static RecordBuilderFactory builderFactory; + + // Large enough to create observable memory pressure + static final int RECORD_COUNT = 100_000; + + // Memory threshold: if buffering all records uses more than this, the test fails. + // 100K records with strings should use >10MB when fully buffered. + // With streaming, memory should stay well under this. + static final long MAX_MEMORY_DELTA_MB = 5; + + @BeforeAll + static void forceManagerInit() { + final ComponentManager manager = ComponentManager.instance(); + if (manager.find(Stream::of).count() == 0) { + manager.addPlugin(new File("target/test-classes").getAbsolutePath()); + } + } + + /** + * Tests memory consumption when a processor emits 100K records in @ElementListener. + * + * FAILS on current code: all 100K records buffered → memory spike > threshold. + * PASSES after iterator implementation: records stream lazily → memory stays low. + */ + @Test + void elementListenerShouldNotBufferAllRecordsInMemory() { + final ComponentManager manager = ComponentManager.instance(); + final Map, Object> servicesMapper = getServicesMapper(manager); + final Jsonb jsonb = (Jsonb) servicesMapper.get(Jsonb.class); + builderFactory = (RecordBuilderFactory) servicesMapper.get(RecordBuilderFactory.class); + + final org.talend.sdk.component.runtime.output.Processor processor = manager + .findProcessor("ProcessorBufferingTest", "heavyEmitter", 1, new HashMap<>()) + .orElseThrow(() -> new IllegalStateException("processor not found")); + JobStateAware.init(processor, new HashMap<>()); + + final AutoChunkProcessor chunkProcessor = new AutoChunkProcessor(RECORD_COUNT + 1, processor); + chunkProcessor.start(); + + final InputsHandler inputsHandler = new InputsHandler(jsonb, servicesMapper); + inputsHandler.addConnection("FLOW", row1Struct.class); + + final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); + outputsHandler.addConnection("FLOW", row1Struct.class); + + final InputFactory inputFactory = inputsHandler.asInputFactory(); + final OutputFactory outputFactory = outputsHandler.asOutputFactory(); + + // Prepare one input record + final Record inputRecord = builderFactory.newRecordBuilder() + .withString("id", "input-1") + .withString("name", "trigger") + .build(); + final RecordConverters.MappingMetaRegistry registry = new RecordConverters.MappingMetaRegistry(); + final row1Struct inputRow = (row1Struct) registry.find(row1Struct.class).newInstance(inputRecord); + inputsHandler.setInputValue("FLOW", inputRow); + + // Force GC to get a clean baseline + System.gc(); + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + final long memoryBefore = usedMemoryMB(); + + // --- Producer emits all records in onElement() --- + chunkProcessor.onElement(inputFactory, outputFactory); + + chunkProcessor.flush(outputFactory); + chunkProcessor.stop(); + + // Measure memory AFTER production, BEFORE drain + final long memoryAfterProduction = usedMemoryMB(); + final long memoryDelta = memoryAfterProduction - memoryBefore; + + System.out.println("=== ProcessorBufferingTest: @ElementListener ==="); + System.out.println("Records emitted: " + RECORD_COUNT); + System.out.println("Memory before onElement(): " + memoryBefore + " MB"); + System.out.println("Memory after onElement(): " + memoryAfterProduction + " MB"); + System.out.println("Memory delta: " + memoryDelta + " MB"); + System.out.println("Threshold: " + MAX_MEMORY_DELTA_MB + " MB"); + + // Drain after onElement — same as Studio generated code + int drainedCount = 0; + while (outputsHandler.hasMoreData()) { + outputsHandler.getValue("FLOW"); + drainedCount++; + } + System.out.println("Records drained: " + drainedCount); + + // ASSERTION: memory delta should be under threshold + // FAILS on current code (all records buffered → large delta) + // PASSES after iterator implementation (lazy streaming → small delta) + assertTrue(memoryDelta <= MAX_MEMORY_DELTA_MB, + "Memory delta should be <= " + MAX_MEMORY_DELTA_MB + " MB (streaming), " + + "but was " + memoryDelta + " MB (all records buffered in memory)"); + } + + /** + * Tests memory consumption when a processor emits 100K records in @AfterGroup. + * + * FAILS on current code: all records buffered in @AfterGroup. + * PASSES after iterator implementation. + */ + @Test + void afterGroupShouldNotBufferAllRecordsInMemory() { + final ComponentManager manager = ComponentManager.instance(); + final Map, Object> servicesMapper = getServicesMapper(manager); + final Jsonb jsonb = (Jsonb) servicesMapper.get(Jsonb.class); + builderFactory = (RecordBuilderFactory) servicesMapper.get(RecordBuilderFactory.class); + + final org.talend.sdk.component.runtime.output.Processor processor = manager + .findProcessor("ProcessorBufferingTest", "heavyAfterGroupEmitter", 1, new HashMap<>()) + .orElseThrow(() -> new IllegalStateException("processor not found")); + JobStateAware.init(processor, new HashMap<>()); + + // chunkSize=5: after 5 inputs, @AfterGroup fires and emits 100K records + final AutoChunkProcessor chunkProcessor = new AutoChunkProcessor(5, processor); + chunkProcessor.start(); + + final InputsHandler inputsHandler = new InputsHandler(jsonb, servicesMapper); + inputsHandler.addConnection("FLOW", row1Struct.class); + + final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); + outputsHandler.addConnection("MAIN", row1Struct.class); + outputsHandler.addConnection("REJECT", row1Struct.class); + + final InputFactory inputFactory = inputsHandler.asInputFactory(); + final OutputFactory outputFactory = outputsHandler.asOutputFactory(); + + final RecordConverters.MappingMetaRegistry registry = new RecordConverters.MappingMetaRegistry(); + + int mainCount = 0; + int rejectCount = 0; + + // Force GC to get a clean baseline + System.gc(); + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + final long memoryBefore = usedMemoryMB(); + + // Feed 5 inputs (triggers @AfterGroup) + for (int i = 0; i < RECORD_COUNT; i++) { + final Record inputRecord = builderFactory.newRecordBuilder() + .withString("id", "input-" + i) + .withString("name", "record-" + i) + .build(); + final row1Struct inputRow = (row1Struct) registry.find(row1Struct.class).newInstance(inputRecord); + inputsHandler.setInputValue("FLOW", inputRow); + + chunkProcessor.onElement(inputFactory, outputFactory); + + // Drain after each onElement — same as Studio generated code + while (outputsHandler.hasMoreData()) { + if (outputsHandler.getValue("MAIN") != null) { + mainCount++; + } + if (outputsHandler.getValue("REJECT") != null) { + rejectCount++; + } + } + } + + chunkProcessor.flush(outputFactory); + chunkProcessor.stop(); + + final long memoryAfterGroup = usedMemoryMB(); + final long memoryDelta = memoryAfterGroup - memoryBefore; + + System.out.println("=== ProcessorBufferingTest: @AfterGroup ==="); + System.out.println("Records emitted: " + RECORD_COUNT); + System.out.println("Memory before: " + memoryBefore + " MB"); + System.out.println("Memory after @AfterGroup: " + memoryAfterGroup + " MB"); + System.out.println("Memory delta: " + memoryDelta + " MB"); + + while (outputsHandler.hasMoreData()) { + if (outputsHandler.getValue("MAIN") != null) { + mainCount++; + } + if (outputsHandler.getValue("REJECT") != null) { + rejectCount++; + } + } + + System.out.println("MAIN drained: " + mainCount); + System.out.println("REJECT drained: " + rejectCount); + + // Verify all records were produced + assertTrue(mainCount + rejectCount > 0, + "Should have produced records in @AfterGroup"); + + // ASSERTION: memory delta should be under threshold + // FAILS on current code (all records buffered → large delta) + // PASSES after iterator implementation (lazy streaming → small delta) + assertTrue(memoryDelta <= MAX_MEMORY_DELTA_MB, + "Memory delta should be <= " + MAX_MEMORY_DELTA_MB + " MB (streaming), " + + "but was " + memoryDelta + " MB (all records buffered in memory)"); + } + + // --- Helpers --- + + private static long usedMemoryMB() { + final Runtime rt = Runtime.getRuntime(); + return (rt.totalMemory() - rt.freeMemory()) / (1024 * 1024); + } + + private Map, Object> getServicesMapper(ComponentManager manager) { + return manager + .findPlugin("test-classes") + .get() + .get(ComponentManager.AllServices.class) + .getServices(); + } + + // --- Test components --- + + /** + * Processor that emits 100K records in @ElementListener via emit(). + * Each record has enough data to create observable memory pressure. + */ + @Processor(name = "heavyEmitter", family = "ProcessorBufferingTest") + public static class HeavyEmitterProcessor implements Serializable { + + @ElementListener + public void onElement(@Input final Record input, + @Output final OutputEmitter output) { + for (int i = 0; i < RECORD_COUNT; i++) { + final Record record = builderFactory.newRecordBuilder() + .withString("id", "out-" + i) + .withString("name", "generated-record-with-some-payload-" + i) + .withString("data", "additional-field-to-increase-memory-footprint-" + i) + .build(); + output.emit(record); + } + } + } + + /** + * Processor that accumulates inputs, then emits 100K records in @AfterGroup. + * Simulates the N:M batch pattern. + */ + @Processor(name = "heavyAfterGroupEmitter", family = "ProcessorBufferingTest") + public static class HeavyAfterGroupEmitterProcessor implements Serializable { + + private int inputCount = 0; + + @BeforeGroup + public void beforeGroup() { + // no-op + } + + @ElementListener + public void onElement(@Input final Record input) { + inputCount++; + } + + @AfterGroup + public void afterGroup(@Output("MAIN") final OutputEmitter main, + @Output("REJECT") final OutputEmitter reject, + @LastGroup final boolean lastGroup) { + if (!lastGroup) { + return; + } + + for (int i = 0; i < RECORD_COUNT; i++) { + final Record record = builderFactory.newRecordBuilder() + .withString("id", "result-" + i) + .withString("name", "processed-record-with-payload-" + i) + .withString("data", "bulk-result-data-field-" + i) + .build(); + if (i % 10 == 0) { + reject.emit(record); + } else { + main.emit(record); + } + } + inputCount = 0; + } + } + + // --- Row struct --- + + @Getter + @ToString + public static class row1Struct implements routines.system.IPersistableRow { + + public String id; + + public String name; + + public String data; + + @Override + public void writeData(final ObjectOutputStream objectOutputStream) { + throw new UnsupportedOperationException("#writeData()"); + } + + @Override + public void readData(final ObjectInputStream objectInputStream) { + throw new UnsupportedOperationException("#readData()"); + } + } +} From 9aecf24a831f908e2612ec6ffd17173510f83b38 Mon Sep 17 00:00:00 2001 From: wwang Date: Mon, 13 Jul 2026 14:01:30 +0800 Subject: [PATCH 02/13] fix(QTDI-2709): do real implement --- .../sdk/component/api/processor/Output.java | 2 + .../api/processor/OutputIterator.java | 55 +++++++++ .../runtime/output/ProcessorImpl.java | 17 ++- .../runtime/visitor/ModelVisitor.java | 11 +- .../component/runtime/di/BaseIOHandler.java | 29 ++++- .../component/runtime/di/OutputsHandler.java | 64 ++++++++-- .../components/ProcessorBufferingTest.java | 112 +++++++++++++----- 7 files changed, 242 insertions(+), 48 deletions(-) create mode 100644 component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java index a6fde9f1e70b2..25748a032d2c1 100644 --- a/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java @@ -26,4 +26,6 @@ public @interface Output { String value() default "__default__"; + + boolean iterator() default false; } diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java new file mode 100644 index 0000000000000..5804dc8e6bf67 --- /dev/null +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java @@ -0,0 +1,55 @@ +/** + * Copyright (C) 2006-2026 Talend Inc. - www.talend.com + * + * Licensed 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.talend.sdk.component.api.processor; + +import java.util.Iterator; + +/** + * Allows a processor to provide a lazy iterator for output records + * instead of pushing them via {@link OutputEmitter#emit(Object)}. + * + *

+ * Used with {@code @Output(iterator = true)} to enable streaming: + * records are produced on-demand during the consumer's drain loop, + * rather than buffered entirely in memory. + * + *

+ * Example usage in a processor: + * + *

+ * {@code
+ * 
+ * @ElementListener
+ * public void process(@Input Record input,
+ *         @Output(iterator = true) OutputIterator output) {
+ *     output.setIterator(myLazyIterator(input));
+ * }
+ * }
+ * 
+ * + * @param the record type + */ +public interface OutputIterator { + + /** + * Sets the lazy iterator that will produce output records on demand. + * Call this once per invocation; the consumer will pull records + * by calling {@link Iterator#hasNext()} and {@link Iterator#next()}. + * + * @param iterator the lazy iterator providing output records + */ + void setIterator(Iterator iterator); +} diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java index 81229696456b1..c69f8b0eef463 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java @@ -56,6 +56,7 @@ import org.talend.sdk.component.api.processor.Input; import org.talend.sdk.component.api.processor.LastGroup; import org.talend.sdk.component.api.processor.Output; +import org.talend.sdk.component.api.processor.OutputIterator; import org.talend.sdk.component.api.service.record.RecordBuilderFactory; import org.talend.sdk.component.runtime.base.Delegated; import org.talend.sdk.component.runtime.base.LifecycleImpl; @@ -150,10 +151,11 @@ public void beforeGroup() { private BiFunction buildProcessParamBuilder(final Parameter parameter) { if (parameter.isAnnotationPresent(Output.class)) { - return (inputs, outputs) -> { - final String name = parameter.getAnnotation(Output.class).value(); - return outputs.create(name); - }; + final Output annotation = parameter.getAnnotation(Output.class); + if (annotation.iterator()) { + return (inputs, outputs) -> (OutputIterator) outputs.create(annotation.value()); + } + return (inputs, outputs) -> outputs.create(annotation.value()); } final Class parameterType = parameter.getType(); @@ -167,8 +169,11 @@ private Function toOutputParamBuilder(final Parameter par if (parameter.isAnnotationPresent(LastGroup.class)) { return false; } - final String name = parameter.getAnnotation(Output.class).value(); - return outputs.create(name); + final Output annotation = parameter.getAnnotation(Output.class); + if (annotation.iterator()) { + return (OutputIterator) outputs.create(annotation.value()); + } + return outputs.create(annotation.value()); }; } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java index 1b52f1b9ff546..8514226d1db9c 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java @@ -46,6 +46,7 @@ import org.talend.sdk.component.api.processor.LastGroup; import org.talend.sdk.component.api.processor.Output; import org.talend.sdk.component.api.processor.OutputEmitter; +import org.talend.sdk.component.api.processor.OutputIterator; import org.talend.sdk.component.api.processor.Processor; import org.talend.sdk.component.api.standalone.DriverRunner; import org.talend.sdk.component.api.standalone.RunAtDriver; @@ -195,7 +196,8 @@ private void validateProcessor(final Class input) { afterGroups.forEach(m -> { final List invalidParams = Stream.of(m.getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { - throw new IllegalArgumentException("@Output parameter must be of type OutputEmitter"); + throw new IllegalArgumentException( + "@Output parameter must be of type OutputEmitter or OutputIterator (with iterator = true)"); } }) .filter(p -> !p.isAnnotationPresent(Output.class)) @@ -243,7 +245,8 @@ private void validateProducer(final Class input, final List afterGrou if (!producers.isEmpty() && Stream.of(producers.get(0).getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { - throw new IllegalArgumentException("@Output parameter must be of type OutputEmitter"); + throw new IllegalArgumentException( + "@Output parameter must be of type OutputEmitter or OutputIterator (with iterator = true)"); } }).filter(p -> !p.isAnnotationPresent(Output.class)).count() < 1) { throw new IllegalArgumentException(input + " doesn't have the input parameter on its producer method"); @@ -254,6 +257,10 @@ private boolean validOutputParam(final Parameter p) { if (!(p.getParameterizedType() instanceof ParameterizedType pt)) { return false; } + final Output annotation = p.getAnnotation(Output.class); + if (annotation != null && annotation.iterator()) { + return OutputIterator.class == pt.getRawType(); + } return OutputEmitter.class == pt.getRawType(); } diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java index 9b7441255171f..4a5f37b4478f2 100644 --- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java +++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java @@ -91,19 +91,32 @@ static class IO { private final Class type; + private Iterator source; + + void setSource(final Iterator source) { + // Close previous source if it implements AutoCloseable + closeSource(); + this.source = source; + } + private void reset() { values.clear(); + closeSource(); + this.source = null; } boolean hasNext() { - return values.size() != 0; + return !values.isEmpty() + || (source != null && source.hasNext()); } T next() { - if (hasNext()) { + if (!values.isEmpty()) { return type.cast(values.poll()); } - + if (source != null && source.hasNext()) { + return type.cast(source.next()); + } return null; } @@ -114,6 +127,16 @@ void add(final T e) { Class getType() { return type; } + + private void closeSource() { + if (source instanceof AutoCloseable) { + try { + ((AutoCloseable) source).close(); + } catch (final Exception e) { + // best effort cleanup + } + } + } } } diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java index 13bbee78e8ce9..3eb9b1d2e3662 100644 --- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java +++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java @@ -15,10 +15,13 @@ */ package org.talend.sdk.component.runtime.di; +import java.util.Iterator; import java.util.Map; import javax.json.bind.Jsonb; +import org.talend.sdk.component.api.processor.OutputEmitter; +import org.talend.sdk.component.api.processor.OutputIterator; import org.talend.sdk.component.api.record.Record; import org.talend.sdk.component.api.record.Schema; import org.talend.sdk.component.runtime.output.OutputFactory; @@ -33,17 +36,9 @@ public OutputsHandler(final Jsonb jsonb, final Map, Object> servicesMap } public OutputFactory asOutputFactory() { - return name -> value -> { + return name -> { final BaseIOHandler.IO ref = connections.get(getActualName(name)); - if (ref != null && value != null) { - if (value instanceof javax.json.JsonValue) { - ref.add(jsonb.fromJson(value.toString(), ref.getType())); - } else if (value instanceof Record record) { - ref.add(registry.find(ref.getType()).newInstance(record)); - } else { - ref.add(jsonb.fromJson(jsonb.toJson(value), ref.getType())); - } - } + return new OutputEmitterWithIterator(ref); }; } @@ -70,4 +65,53 @@ public OutputFactory asOutputFactoryForGuessSchema() { }; } + private Object convert(final Object value, final BaseIOHandler.IO ref) { + if (value == null) { + return null; + } else if (value instanceof javax.json.JsonValue) { + return jsonb.fromJson(value.toString(), ref.getType()); + } else if (value instanceof Record record) { + return registry.find(ref.getType()).newInstance(record); + } else { + return jsonb.fromJson(jsonb.toJson(value), ref.getType()); + } + } + + /** + * Internal class implementing both OutputEmitter and OutputIterator. + * Allows ProcessorImpl to cast to OutputIterator when iterator mode is active. + */ + private class OutputEmitterWithIterator implements OutputEmitter, OutputIterator { + + private final BaseIOHandler.IO ref; + + OutputEmitterWithIterator(final BaseIOHandler.IO ref) { + this.ref = ref; + } + + @Override + public void emit(final Object value) { + if (ref != null && value != null) { + ref.add(convert(value, ref)); + } + } + + @Override + public void setIterator(final Iterator iterator) { + if (ref != null) { + ref.setSource(new Iterator() { + + @Override + public boolean hasNext() { + return iterator.hasNext(); + } + + @Override + public Object next() { + return convert(iterator.next(), ref); + } + }); + } + } + } } diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java index 299b97edfc068..062f245d65072 100644 --- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java @@ -35,7 +35,7 @@ import org.talend.sdk.component.api.processor.Input; import org.talend.sdk.component.api.processor.LastGroup; import org.talend.sdk.component.api.processor.Output; -import org.talend.sdk.component.api.processor.OutputEmitter; +import org.talend.sdk.component.api.processor.OutputIterator; import org.talend.sdk.component.api.processor.Processor; import org.talend.sdk.component.api.record.Record; import org.talend.sdk.component.api.service.record.RecordBuilderFactory; @@ -233,6 +233,14 @@ void afterGroupShouldNotBufferAllRecordsInMemory() { chunkProcessor.flush(outputFactory); chunkProcessor.stop(); + // GC to reclaim input record objects from the loop + System.gc(); + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + final long memoryAfterGroup = usedMemoryMB(); final long memoryDelta = memoryAfterGroup - memoryBefore; @@ -284,29 +292,41 @@ private Map, Object> getServicesMapper(ComponentManager manager) { // --- Test components --- /** - * Processor that emits 100K records in @ElementListener via emit(). - * Each record has enough data to create observable memory pressure. + * Processor that provides a lazy iterator in @ElementListener. + * Records are produced on-demand during the drain loop. */ @Processor(name = "heavyEmitter", family = "ProcessorBufferingTest") public static class HeavyEmitterProcessor implements Serializable { @ElementListener public void onElement(@Input final Record input, - @Output final OutputEmitter output) { - for (int i = 0; i < RECORD_COUNT; i++) { - final Record record = builderFactory.newRecordBuilder() - .withString("id", "out-" + i) - .withString("name", "generated-record-with-some-payload-" + i) - .withString("data", "additional-field-to-increase-memory-footprint-" + i) - .build(); - output.emit(record); - } + @Output(iterator = true) final OutputIterator output) { + output.setIterator(new java.util.Iterator<>() { + + private int index = 0; + + @Override + public boolean hasNext() { + return index < RECORD_COUNT; + } + + @Override + public Record next() { + final Record record = builderFactory.newRecordBuilder() + .withString("id", "out-" + index) + .withString("name", "generated-record-with-some-payload-" + index) + .withString("data", "additional-field-to-increase-memory-footprint-" + index) + .build(); + index++; + return record; + } + }); } } /** - * Processor that accumulates inputs, then emits 100K records in @AfterGroup. - * Simulates the N:M batch pattern. + * Processor that accumulates inputs, then sets lazy iterators on MAIN and REJECT + * in @AfterGroup. Simulates the N:M batch pattern. */ @Processor(name = "heavyAfterGroupEmitter", family = "ProcessorBufferingTest") public static class HeavyAfterGroupEmitterProcessor implements Serializable { @@ -324,25 +344,63 @@ public void onElement(@Input final Record input) { } @AfterGroup - public void afterGroup(@Output("MAIN") final OutputEmitter main, - @Output("REJECT") final OutputEmitter reject, + public void afterGroup(@Output(value = "MAIN", iterator = true) final OutputIterator main, + @Output(value = "REJECT", iterator = true) final OutputIterator reject, @LastGroup final boolean lastGroup) { if (!lastGroup) { return; } - for (int i = 0; i < RECORD_COUNT; i++) { - final Record record = builderFactory.newRecordBuilder() - .withString("id", "result-" + i) - .withString("name", "processed-record-with-payload-" + i) - .withString("data", "bulk-result-data-field-" + i) - .build(); - if (i % 10 == 0) { - reject.emit(record); - } else { - main.emit(record); + // MAIN gets records where index % 10 != 0 + main.setIterator(new java.util.Iterator<>() { + + private int index = 0; + + @Override + public boolean hasNext() { + while (index < RECORD_COUNT && index % 10 == 0) { + index++; + } + return index < RECORD_COUNT; } - } + + @Override + public Record next() { + final Record record = builderFactory.newRecordBuilder() + .withString("id", "result-" + index) + .withString("name", "processed-record-with-payload-" + index) + .withString("data", "bulk-result-data-field-" + index) + .build(); + index++; + return record; + } + }); + + // REJECT gets records where index % 10 == 0 + reject.setIterator(new java.util.Iterator<>() { + + private int index = 0; + + @Override + public boolean hasNext() { + while (index < RECORD_COUNT && index % 10 != 0) { + index++; + } + return index < RECORD_COUNT; + } + + @Override + public Record next() { + final Record record = builderFactory.newRecordBuilder() + .withString("id", "reject-" + index) + .withString("name", "rejected-record-" + index) + .withString("data", "reject-data-" + index) + .build(); + index++; + return record; + } + }); + inputCount = 0; } } From b33e762cd69ab73872c5f3721b6d4c5a7f739fc4 Mon Sep 17 00:00:00 2001 From: wwang Date: Mon, 13 Jul 2026 14:32:55 +0800 Subject: [PATCH 03/13] fix(QTDI-2709): remove verbose iterator attribute --- .../sdk/component/api/processor/Output.java | 2 -- .../component/api/processor/OutputIterator.java | 4 ++-- .../component/runtime/output/ProcessorImpl.java | 16 ++++++++-------- .../component/runtime/visitor/ModelVisitor.java | 10 +++------- .../beam/components/ProcessorBufferingTest.java | 6 +++--- .../component/tools/ComponentValidatorTest.java | 2 +- 6 files changed, 17 insertions(+), 23 deletions(-) diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java index 25748a032d2c1..a6fde9f1e70b2 100644 --- a/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/Output.java @@ -26,6 +26,4 @@ public @interface Output { String value() default "__default__"; - - boolean iterator() default false; } diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java index 5804dc8e6bf67..af39bd436b025 100644 --- a/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java @@ -22,7 +22,7 @@ * instead of pushing them via {@link OutputEmitter#emit(Object)}. * *

- * Used with {@code @Output(iterator = true)} to enable streaming: + * Used with {@code @Output} on an {@code OutputIterator} parameter to enable streaming: * records are produced on-demand during the consumer's drain loop, * rather than buffered entirely in memory. * @@ -34,7 +34,7 @@ * * @ElementListener * public void process(@Input Record input, - * @Output(iterator = true) OutputIterator output) { + * @Output OutputIterator output) { * output.setIterator(myLazyIterator(input)); * } * } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java index c69f8b0eef463..708695fc5d0e4 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java @@ -151,11 +151,11 @@ public void beforeGroup() { private BiFunction buildProcessParamBuilder(final Parameter parameter) { if (parameter.isAnnotationPresent(Output.class)) { - final Output annotation = parameter.getAnnotation(Output.class); - if (annotation.iterator()) { - return (inputs, outputs) -> (OutputIterator) outputs.create(annotation.value()); + final String name = parameter.getAnnotation(Output.class).value(); + if (OutputIterator.class == parameter.getType()) { + return (inputs, outputs) -> (OutputIterator) outputs.create(name); } - return (inputs, outputs) -> outputs.create(annotation.value()); + return (inputs, outputs) -> outputs.create(name); } final Class parameterType = parameter.getType(); @@ -169,11 +169,11 @@ private Function toOutputParamBuilder(final Parameter par if (parameter.isAnnotationPresent(LastGroup.class)) { return false; } - final Output annotation = parameter.getAnnotation(Output.class); - if (annotation.iterator()) { - return (OutputIterator) outputs.create(annotation.value()); + final String name = parameter.getAnnotation(Output.class).value(); + if (OutputIterator.class == parameter.getType()) { + return (OutputIterator) outputs.create(name); } - return outputs.create(annotation.value()); + return outputs.create(name); }; } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java index 8514226d1db9c..06d9dadadaee2 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java @@ -197,7 +197,7 @@ private void validateProcessor(final Class input) { final List invalidParams = Stream.of(m.getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { throw new IllegalArgumentException( - "@Output parameter must be of type OutputEmitter or OutputIterator (with iterator = true)"); + "@Output parameter must be of type OutputEmitter or OutputIterator"); } }) .filter(p -> !p.isAnnotationPresent(Output.class)) @@ -246,7 +246,7 @@ private void validateProducer(final Class input, final List afterGrou if (!producers.isEmpty() && Stream.of(producers.get(0).getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { throw new IllegalArgumentException( - "@Output parameter must be of type OutputEmitter or OutputIterator (with iterator = true)"); + "@Output parameter must be of type OutputEmitter or OutputIterator"); } }).filter(p -> !p.isAnnotationPresent(Output.class)).count() < 1) { throw new IllegalArgumentException(input + " doesn't have the input parameter on its producer method"); @@ -257,11 +257,7 @@ private boolean validOutputParam(final Parameter p) { if (!(p.getParameterizedType() instanceof ParameterizedType pt)) { return false; } - final Output annotation = p.getAnnotation(Output.class); - if (annotation != null && annotation.iterator()) { - return OutputIterator.class == pt.getRawType(); - } - return OutputEmitter.class == pt.getRawType(); + return OutputEmitter.class == pt.getRawType() || OutputIterator.class == pt.getRawType(); } private Stream> getPartitionMapperMethods(final boolean infinite) { diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java index 062f245d65072..952bb8618c754 100644 --- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java @@ -300,7 +300,7 @@ public static class HeavyEmitterProcessor implements Serializable { @ElementListener public void onElement(@Input final Record input, - @Output(iterator = true) final OutputIterator output) { + @Output final OutputIterator output) { output.setIterator(new java.util.Iterator<>() { private int index = 0; @@ -344,8 +344,8 @@ public void onElement(@Input final Record input) { } @AfterGroup - public void afterGroup(@Output(value = "MAIN", iterator = true) final OutputIterator main, - @Output(value = "REJECT", iterator = true) final OutputIterator reject, + public void afterGroup(@Output("MAIN") final OutputIterator main, + @Output("REJECT") final OutputIterator reject, @LastGroup final boolean lastGroup) { if (!lastGroup) { return; diff --git a/component-tools/src/test/java/org/talend/sdk/component/tools/ComponentValidatorTest.java b/component-tools/src/test/java/org/talend/sdk/component/tools/ComponentValidatorTest.java index d4aa89522267a..65bdc0ea5e49b 100755 --- a/component-tools/src/test/java/org/talend/sdk/component/tools/ComponentValidatorTest.java +++ b/component-tools/src/test/java/org/talend/sdk/component/tools/ComponentValidatorTest.java @@ -571,7 +571,7 @@ void testFailureAfterGroup(final ExceptionSpec expectedException) { expectedException .expectMessage( """ - - @Output parameter must be of type OutputEmitter + - @Output parameter must be of type OutputEmitter or OutputIterator - Parameter of AfterGroup method need to be annotated with Output - class org.talend.test.failure.aftergroup.MyComponent5 must have a single @AfterGroup method with @LastGroup parameter"""); } From 87deb2d6a4c770288667cdfedc3b3e334ebd1832 Mon Sep 17 00:00:00 2001 From: wwang Date: Mon, 13 Jul 2026 15:22:22 +0800 Subject: [PATCH 04/13] doc(QTDI-2709): update OutputIterator documentation to clarify runtime support limitations --- .../talend/sdk/component/api/processor/OutputIterator.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java index af39bd436b025..cfd62d32ebfe6 100644 --- a/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/OutputIterator.java @@ -27,6 +27,12 @@ * rather than buffered entirely in memory. * *

+ * Important: This interface is supported only in the Studio DI runtime. + * It is not supported in Beam-based runners (Cloud) where processors typically + * do not use TCK processor patterns. Using {@code OutputIterator} in a Beam + * pipeline will result in a {@code ClassCastException} at runtime. + * + *

* Example usage in a processor: * *


From 9e8138b85a1fd71eb291340d976b892b810b4241 Mon Sep 17 00:00:00 2001
From: wwang 
Date: Mon, 13 Jul 2026 17:20:18 +0800
Subject: [PATCH 05/13] fix(QTDI-2709): refactor variable names for clarity in
 OutputsHandler and ProcessorBufferingTest

---
 .../sdk/component/runtime/di/OutputsHandler.java   | 14 +++++++-------
 .../di/beam/components/ProcessorBufferingTest.java | 12 ++++++------
 2 files changed, 13 insertions(+), 13 deletions(-)

diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java
index 3eb9b1d2e3662..0ee35afe590c3 100644
--- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java
+++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java
@@ -54,8 +54,8 @@ public OutputFactory asOutputFactoryForGuessSchema() {
             if (ref != null && value != null) {
                 if (value instanceof javax.json.JsonValue) {
                     ref.add(jsonb.fromJson(value.toString(), ref.getType()));
-                } else if (value instanceof Record record) {
-                    ref.add(record.getSchema());
+                } else if (value instanceof Record rec) {
+                    ref.add(rec.getSchema());
                 } else if (value instanceof Schema) {
                     ref.add(value);
                 } else {
@@ -70,8 +70,8 @@ private Object convert(final Object value, final BaseIOHandler.IO ref) {
             return null;
         } else if (value instanceof javax.json.JsonValue) {
             return jsonb.fromJson(value.toString(), ref.getType());
-        } else if (value instanceof Record record) {
-            return registry.find(ref.getType()).newInstance(record);
+        } else if (value instanceof Record rec) {
+            return registry.find(ref.getType()).newInstance(rec);
         } else {
             return jsonb.fromJson(jsonb.toJson(value), ref.getType());
         }
@@ -81,7 +81,7 @@ private Object convert(final Object value, final BaseIOHandler.IO ref) {
      * Internal class implementing both OutputEmitter and OutputIterator.
      * Allows ProcessorImpl to cast to OutputIterator when iterator mode is active.
      */
-    private class OutputEmitterWithIterator implements OutputEmitter, OutputIterator {
+    private class OutputEmitterWithIterator implements OutputEmitter, OutputIterator {
 
         private final BaseIOHandler.IO ref;
 
@@ -90,14 +90,14 @@ private class OutputEmitterWithIterator implements OutputEmitter, OutputIterator
         }
 
         @Override
-        public void emit(final Object value) {
+        public void emit(final T value) {
             if (ref != null && value != null) {
                 ref.add(convert(value, ref));
             }
         }
 
         @Override
-        public void setIterator(final Iterator iterator) {
+        public void setIterator(final Iterator iterator) {
             if (ref != null) {
                 ref.setSource(new Iterator() {
 
diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java
index 952bb8618c754..906e1004e274b 100644
--- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java
+++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java
@@ -312,13 +312,13 @@ public boolean hasNext() {
 
                 @Override
                 public Record next() {
-                    final Record record = builderFactory.newRecordBuilder()
+                    final Record rec = builderFactory.newRecordBuilder()
                             .withString("id", "out-" + index)
                             .withString("name", "generated-record-with-some-payload-" + index)
                             .withString("data", "additional-field-to-increase-memory-footprint-" + index)
                             .build();
                     index++;
-                    return record;
+                    return rec;
                 }
             });
         }
@@ -366,13 +366,13 @@ public boolean hasNext() {
 
                 @Override
                 public Record next() {
-                    final Record record = builderFactory.newRecordBuilder()
+                    final Record rec = builderFactory.newRecordBuilder()
                             .withString("id", "result-" + index)
                             .withString("name", "processed-record-with-payload-" + index)
                             .withString("data", "bulk-result-data-field-" + index)
                             .build();
                     index++;
-                    return record;
+                    return rec;
                 }
             });
 
@@ -391,13 +391,13 @@ public boolean hasNext() {
 
                 @Override
                 public Record next() {
-                    final Record record = builderFactory.newRecordBuilder()
+                    final Record rec = builderFactory.newRecordBuilder()
                             .withString("id", "reject-" + index)
                             .withString("name", "rejected-record-" + index)
                             .withString("data", "reject-data-" + index)
                             .build();
                     index++;
-                    return record;
+                    return rec;
                 }
             });
 

From 0c7f73c5c1958c8782334d38e9700db9cf44fea8 Mon Sep 17 00:00:00 2001
From: wwang 
Date: Tue, 14 Jul 2026 11:19:56 +0800
Subject: [PATCH 06/13] fix(QTDI-2709): implement OutputIterator support in
 OutputFactory and OutputsHandler

---
 .../runtime/output/OutputFactory.java         | 15 ++++
 .../runtime/output/ProcessorImpl.java         |  4 +-
 .../component/runtime/di/BaseIOHandler.java   | 17 ++++-
 .../component/runtime/di/OutputsHandler.java  | 73 +++++++++----------
 4 files changed, 66 insertions(+), 43 deletions(-)

diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java
index c3f8b14bf5fc1..b255c29924b25 100644
--- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java
+++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java
@@ -16,8 +16,23 @@
 package org.talend.sdk.component.runtime.output;
 
 import org.talend.sdk.component.api.processor.OutputEmitter;
+import org.talend.sdk.component.api.processor.OutputIterator;
 
 public interface OutputFactory {
 
     OutputEmitter create(String name);
+
+    /**
+     * Creates an {@link OutputIterator} for the given output branch name.
+     * Supported only in the Studio DI runtime. Other runtimes throw
+     * {@link UnsupportedOperationException} by default.
+     *
+     * @param name the output branch name
+     * @return an OutputIterator for lazy record streaming
+     * @throws UnsupportedOperationException if the runtime does not support iterator mode
+     */
+    default OutputIterator createIterator(String name) {
+        throw new UnsupportedOperationException(
+                "OutputIterator is only supported in the Studio DI runtime");
+    }
 }
diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java
index 708695fc5d0e4..38aab8af13635 100644
--- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java
+++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java
@@ -153,7 +153,7 @@ private BiFunction buildProcessParamBuilder
         if (parameter.isAnnotationPresent(Output.class)) {
             final String name = parameter.getAnnotation(Output.class).value();
             if (OutputIterator.class == parameter.getType()) {
-                return (inputs, outputs) -> (OutputIterator) outputs.create(name);
+                return (inputs, outputs) -> outputs.createIterator(name);
             }
             return (inputs, outputs) -> outputs.create(name);
         }
@@ -171,7 +171,7 @@ private Function toOutputParamBuilder(final Parameter par
             }
             final String name = parameter.getAnnotation(Output.class).value();
             if (OutputIterator.class == parameter.getType()) {
-                return (OutputIterator) outputs.create(name);
+                return outputs.createIterator(name);
             }
             return outputs.create(name);
         };
diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java
index 4a5f37b4478f2..a25adff74677e 100644
--- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java
+++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/BaseIOHandler.java
@@ -29,7 +29,9 @@
 import org.talend.sdk.component.runtime.record.RecordConverters;
 
 import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
 
+@Slf4j
 public abstract class BaseIOHandler {
 
     protected final Jsonb jsonb;
@@ -84,6 +86,19 @@ protected String getActualName(final String name) {
         return "__default__".equals(name) ? "FLOW" : name;
     }
 
+    /**
+     * Represents a single output connection's data holder.
+     * Supports two modes (mutually exclusive per invocation):
+     * 
    + *
  • Push mode (default): records are added via {@link #add(Object)} into the internal queue. + * Used by {@code OutputEmitter.emit()}.
  • + *
  • Pull mode (iterator): a lazy {@link Iterator} source is set via {@link #setSource(Iterator)}. + * Records are produced on-demand during the drain loop ({@link #hasNext()}/{@link #next()}). + * Used by {@code OutputIterator.setIterator()} in the Studio DI runtime.
  • + *
+ * Note: The {@code source} field is only used by {@code OutputsHandler}; {@code InputsHandler} + * uses only the queue-based push mode. + */ @RequiredArgsConstructor static class IO { @@ -133,7 +148,7 @@ private void closeSource() { try { ((AutoCloseable) source).close(); } catch (final Exception e) { - // best effort cleanup + log.debug("Failed to close iterator source", e); } } } diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java index 0ee35afe590c3..ed5e949f25d53 100644 --- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java +++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java @@ -36,9 +36,39 @@ public OutputsHandler(final Jsonb jsonb, final Map, Object> servicesMap } public OutputFactory asOutputFactory() { - return name -> { - final BaseIOHandler.IO ref = connections.get(getActualName(name)); - return new OutputEmitterWithIterator(ref); + return new OutputFactory() { + + @Override + public OutputEmitter create(final String name) { + final BaseIOHandler.IO ref = connections.get(getActualName(name)); + return value -> { + if (ref != null && value != null) { + ref.add(convert(value, ref)); + } + }; + } + + @Override + public OutputIterator createIterator(final String name) { + final BaseIOHandler.IO ref = connections.get(getActualName(name)); + return iterator -> { + if (ref == null) { + return; + } + ref.setSource(new Iterator() { + + @Override + public boolean hasNext() { + return iterator.hasNext(); + } + + @Override + public Object next() { + return convert(iterator.next(), ref); + } + }); + }; + } }; } @@ -77,41 +107,4 @@ private Object convert(final Object value, final BaseIOHandler.IO ref) { } } - /** - * Internal class implementing both OutputEmitter and OutputIterator. - * Allows ProcessorImpl to cast to OutputIterator when iterator mode is active. - */ - private class OutputEmitterWithIterator implements OutputEmitter, OutputIterator { - - private final BaseIOHandler.IO ref; - - OutputEmitterWithIterator(final BaseIOHandler.IO ref) { - this.ref = ref; - } - - @Override - public void emit(final T value) { - if (ref != null && value != null) { - ref.add(convert(value, ref)); - } - } - - @Override - public void setIterator(final Iterator iterator) { - if (ref != null) { - ref.setSource(new Iterator() { - - @Override - public boolean hasNext() { - return iterator.hasNext(); - } - - @Override - public Object next() { - return convert(iterator.next(), ref); - } - }); - } - } - } } From cac943b5decf2f27027baff656988c41fd7bfb06 Mon Sep 17 00:00:00 2001 From: wwang Date: Tue, 14 Jul 2026 11:43:04 +0800 Subject: [PATCH 07/13] fix(QTDI-2709): rename ProcessorBufferingTest package for consistency in component structure --- .../di/{beam/components => studio}/ProcessorBufferingTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) rename component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/{beam/components => studio}/ProcessorBufferingTest.java (99%) diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java similarity index 99% rename from component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java rename to component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java index 906e1004e274b..44a118b272776 100644 --- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/beam/components/ProcessorBufferingTest.java +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.talend.sdk.component.runtime.di.beam.components; +package org.talend.sdk.component.runtime.di.studio; import static org.junit.jupiter.api.Assertions.assertTrue; From 25264bc0c967543d04c3337ad25b254ac0e3e870 Mon Sep 17 00:00:00 2001 From: wwang Date: Thu, 16 Jul 2026 10:11:24 +0800 Subject: [PATCH 08/13] add streaming api --- .../sdk/component/runtime/output/Processor.java | 4 ++++ .../sdk/component/runtime/output/ProcessorImpl.java | 13 +++++++++++++ .../runtime/manager/chain/AutoChunkProcessor.java | 4 ++++ 3 files changed, 21 insertions(+) diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java index 3638e272576c3..66387b0ec9a64 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java @@ -34,5 +34,9 @@ default boolean isLastGroupUsed() { return false; } + default boolean isStreamingMode() { + return false; + } + void onNext(InputFactory input, OutputFactory output); } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java index 38aab8af13635..8226cac9b967b 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java @@ -95,6 +95,8 @@ public class ProcessorImpl extends LifecycleImpl implements Processor, Delegated private transient Collection records; + private transient Boolean streamingMode; + private Map internalConfiguration; private RecordConverters.MappingMetaRegistry mappings; @@ -261,6 +263,17 @@ public void afterGroup(final OutputFactory output) { } } + @Override + public boolean isStreamingMode() { + if (streamingMode == null) { + streamingMode = Stream + .concat(findMethods(ElementListener.class), findMethods(AfterGroup.class)) + .flatMap(m -> Stream.of(m.getParameters())) + .anyMatch(p -> p.isAnnotationPresent(Output.class) && OutputIterator.class == p.getType()); + } + return streamingMode; + } + @Override public boolean isLastGroupUsed() { AtomicReference hasLastGroup = new AtomicReference<>(false); diff --git a/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java b/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java index e1ec2c379bf46..d2c1498adef08 100644 --- a/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java +++ b/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java @@ -46,6 +46,10 @@ public void onElement(final InputFactory ins, final OutputFactory outs) { } } + public boolean isStreamingMode() { + return processor.isStreamingMode(); + } + public void flush(final OutputFactory outs) { if (processedItemCount > 0) { processor.afterGroup(outs); From 1788396938e3105f5b039b1f1000a60b5c44a920 Mon Sep 17 00:00:00 2001 From: wwang Date: Thu, 16 Jul 2026 10:13:21 +0800 Subject: [PATCH 09/13] improve test --- .../di/studio/ProcessorBufferingTest.java | 104 ++++++++++++------ 1 file changed, 69 insertions(+), 35 deletions(-) diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java index 44a118b272776..939d79b33a4ee 100644 --- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java @@ -28,7 +28,8 @@ import javax.json.bind.Jsonb; import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import org.talend.sdk.component.api.processor.AfterGroup; import org.talend.sdk.component.api.processor.BeforeGroup; import org.talend.sdk.component.api.processor.ElementListener; @@ -73,6 +74,12 @@ class ProcessorBufferingTest { // With streaming, memory should stay well under this. static final long MAX_MEMORY_DELTA_MB = 5; + enum OutputMode { + NO_OUTPUT, + ONE_OUTPUT, + TWO_OUTPUT + } + @BeforeAll static void forceManagerInit() { final ComponentManager manager = ComponentManager.instance(); @@ -87,8 +94,9 @@ static void forceManagerInit() { * FAILS on current code: all 100K records buffered → memory spike > threshold. * PASSES after iterator implementation: records stream lazily → memory stays low. */ - @Test - void elementListenerShouldNotBufferAllRecordsInMemory() { + @ParameterizedTest + @EnumSource(OutputMode.class) + void elementListenerShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { final ComponentManager manager = ComponentManager.instance(); final Map, Object> servicesMapper = getServicesMapper(manager); final Jsonb jsonb = (Jsonb) servicesMapper.get(Jsonb.class); @@ -106,7 +114,9 @@ void elementListenerShouldNotBufferAllRecordsInMemory() { inputsHandler.addConnection("FLOW", row1Struct.class); final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); - outputsHandler.addConnection("FLOW", row1Struct.class); + if(outputMode != OutputMode.NO_OUTPUT) { + outputsHandler.addConnection("FLOW", row1Struct.class); + } final InputFactory inputFactory = inputsHandler.asInputFactory(); final OutputFactory outputFactory = outputsHandler.asOutputFactory(); @@ -133,27 +143,37 @@ void elementListenerShouldNotBufferAllRecordsInMemory() { chunkProcessor.onElement(inputFactory, outputFactory); chunkProcessor.flush(outputFactory); - chunkProcessor.stop(); - // Measure memory AFTER production, BEFORE drain + // Measure memory AFTER production, BEFORE drain — valid for all output modes final long memoryAfterProduction = usedMemoryMB(); final long memoryDelta = memoryAfterProduction - memoryBefore; System.out.println("=== ProcessorBufferingTest: @ElementListener ==="); + System.out.println("Output mode: " + outputMode); System.out.println("Records emitted: " + RECORD_COUNT); System.out.println("Memory before onElement(): " + memoryBefore + " MB"); System.out.println("Memory after onElement(): " + memoryAfterProduction + " MB"); System.out.println("Memory delta: " + memoryDelta + " MB"); System.out.println("Threshold: " + MAX_MEMORY_DELTA_MB + " MB"); - // Drain after onElement — same as Studio generated code + // Drain — same as Studio generated code (inside + outside loop, skipped for NO_OUTPUT) int drainedCount = 0; - while (outputsHandler.hasMoreData()) { - outputsHandler.getValue("FLOW"); - drainedCount++; + if (outputMode != OutputMode.NO_OUTPUT) { + // sim inside of input loop + while (outputsHandler.hasMoreData()) { + outputsHandler.getValue("FLOW"); + drainedCount++; + } + // sim outside of input loop + while (outputsHandler.hasMoreData()) { + outputsHandler.getValue("FLOW"); + drainedCount++; + } } System.out.println("Records drained: " + drainedCount); + chunkProcessor.stop(); + // ASSERTION: memory delta should be under threshold // FAILS on current code (all records buffered → large delta) // PASSES after iterator implementation (lazy streaming → small delta) @@ -168,8 +188,9 @@ void elementListenerShouldNotBufferAllRecordsInMemory() { * FAILS on current code: all records buffered in @AfterGroup. * PASSES after iterator implementation. */ - @Test - void afterGroupShouldNotBufferAllRecordsInMemory() { + @ParameterizedTest + @EnumSource(OutputMode.class) + void afterGroupShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { final ComponentManager manager = ComponentManager.instance(); final Map, Object> servicesMapper = getServicesMapper(manager); final Jsonb jsonb = (Jsonb) servicesMapper.get(Jsonb.class); @@ -188,8 +209,12 @@ void afterGroupShouldNotBufferAllRecordsInMemory() { inputsHandler.addConnection("FLOW", row1Struct.class); final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); - outputsHandler.addConnection("MAIN", row1Struct.class); - outputsHandler.addConnection("REJECT", row1Struct.class); + if(outputMode != OutputMode.NO_OUTPUT) { + outputsHandler.addConnection("MAIN", row1Struct.class); + if(outputMode == OutputMode.TWO_OUTPUT) { + outputsHandler.addConnection("REJECT", row1Struct.class); + } + } final InputFactory inputFactory = inputsHandler.asInputFactory(); final OutputFactory outputFactory = outputsHandler.asOutputFactory(); @@ -219,19 +244,22 @@ void afterGroupShouldNotBufferAllRecordsInMemory() { chunkProcessor.onElement(inputFactory, outputFactory); - // Drain after each onElement — same as Studio generated code - while (outputsHandler.hasMoreData()) { - if (outputsHandler.getValue("MAIN") != null) { - mainCount++; - } - if (outputsHandler.getValue("REJECT") != null) { - rejectCount++; + if (outputMode != OutputMode.NO_OUTPUT) { + // Drain after each onElement — same as Studio generated code + while (outputsHandler.hasMoreData()) { + if (outputsHandler.getValue("MAIN") != null) { + mainCount++; + } + if (outputMode == OutputMode.TWO_OUTPUT) { + if (outputsHandler.getValue("REJECT") != null) { + rejectCount++; + } + } } } } chunkProcessor.flush(outputFactory); - chunkProcessor.stop(); // GC to reclaim input record objects from the loop System.gc(); @@ -250,21 +278,27 @@ void afterGroupShouldNotBufferAllRecordsInMemory() { System.out.println("Memory after @AfterGroup: " + memoryAfterGroup + " MB"); System.out.println("Memory delta: " + memoryDelta + " MB"); - while (outputsHandler.hasMoreData()) { - if (outputsHandler.getValue("MAIN") != null) { - mainCount++; - } - if (outputsHandler.getValue("REJECT") != null) { - rejectCount++; + if (outputMode != OutputMode.NO_OUTPUT) { + while (outputsHandler.hasMoreData()) { + if (outputsHandler.getValue("MAIN") != null) { + mainCount++; + } + if (outputMode == OutputMode.TWO_OUTPUT) { + if (outputsHandler.getValue("REJECT") != null) { + rejectCount++; + } + } } + + // Verify all records were produced + assertTrue(mainCount + rejectCount > 0, + "Should have produced records in @AfterGroup"); } System.out.println("MAIN drained: " + mainCount); System.out.println("REJECT drained: " + rejectCount); - // Verify all records were produced - assertTrue(mainCount + rejectCount > 0, - "Should have produced records in @AfterGroup"); + chunkProcessor.stop(); // ASSERTION: memory delta should be under threshold // FAILS on current code (all records buffered → large delta) @@ -351,14 +385,14 @@ public void afterGroup(@Output("MAIN") final OutputIterator main, return; } - // MAIN gets records where index % 10 != 0 + // MAIN gets records where index % 2 != 0 main.setIterator(new java.util.Iterator<>() { private int index = 0; @Override public boolean hasNext() { - while (index < RECORD_COUNT && index % 10 == 0) { + while (index < RECORD_COUNT && index % 2 == 0) { index++; } return index < RECORD_COUNT; @@ -376,14 +410,14 @@ public Record next() { } }); - // REJECT gets records where index % 10 == 0 + // REJECT gets records where index % 2 == 0 reject.setIterator(new java.util.Iterator<>() { private int index = 0; @Override public boolean hasNext() { - while (index < RECORD_COUNT && index % 10 != 0) { + while (index < RECORD_COUNT && index % 2 != 0) { index++; } return index < RECORD_COUNT; From 283bd80331e69a2ccd7395dff04fbfd5d3a0a6a9 Mon Sep 17 00:00:00 2001 From: wwang Date: Thu, 16 Jul 2026 10:13:30 +0800 Subject: [PATCH 10/13] Revert "add streaming api" This reverts commit 25264bc0c967543d04c3337ad25b254ac0e3e870. --- .../sdk/component/runtime/output/Processor.java | 4 ---- .../sdk/component/runtime/output/ProcessorImpl.java | 13 ------------- .../runtime/manager/chain/AutoChunkProcessor.java | 4 ---- 3 files changed, 21 deletions(-) diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java index 66387b0ec9a64..3638e272576c3 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/Processor.java @@ -34,9 +34,5 @@ default boolean isLastGroupUsed() { return false; } - default boolean isStreamingMode() { - return false; - } - void onNext(InputFactory input, OutputFactory output); } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java index 8226cac9b967b..38aab8af13635 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java @@ -95,8 +95,6 @@ public class ProcessorImpl extends LifecycleImpl implements Processor, Delegated private transient Collection records; - private transient Boolean streamingMode; - private Map internalConfiguration; private RecordConverters.MappingMetaRegistry mappings; @@ -263,17 +261,6 @@ public void afterGroup(final OutputFactory output) { } } - @Override - public boolean isStreamingMode() { - if (streamingMode == null) { - streamingMode = Stream - .concat(findMethods(ElementListener.class), findMethods(AfterGroup.class)) - .flatMap(m -> Stream.of(m.getParameters())) - .anyMatch(p -> p.isAnnotationPresent(Output.class) && OutputIterator.class == p.getType()); - } - return streamingMode; - } - @Override public boolean isLastGroupUsed() { AtomicReference hasLastGroup = new AtomicReference<>(false); diff --git a/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java b/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java index d2c1498adef08..e1ec2c379bf46 100644 --- a/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java +++ b/component-runtime-manager/src/main/java/org/talend/sdk/component/runtime/manager/chain/AutoChunkProcessor.java @@ -46,10 +46,6 @@ public void onElement(final InputFactory ins, final OutputFactory outs) { } } - public boolean isStreamingMode() { - return processor.isStreamingMode(); - } - public void flush(final OutputFactory outs) { if (processedItemCount > 0) { processor.afterGroup(outs); From d7273b4a55f50b5f16beffd2d311c3292e534d90 Mon Sep 17 00:00:00 2001 From: wwang Date: Thu, 16 Jul 2026 10:29:44 +0800 Subject: [PATCH 11/13] improve test --- .../di/studio/ProcessorBufferingTest.java | 59 ++++++++++--------- 1 file changed, 30 insertions(+), 29 deletions(-) diff --git a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java index 939d79b33a4ee..0bf158f27c36b 100644 --- a/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java +++ b/component-studio/component-runtime-di/src/test/java/org/talend/sdk/component/runtime/di/studio/ProcessorBufferingTest.java @@ -114,7 +114,7 @@ void elementListenerShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { inputsHandler.addConnection("FLOW", row1Struct.class); final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); - if(outputMode != OutputMode.NO_OUTPUT) { + if (outputMode != OutputMode.NO_OUTPUT) { outputsHandler.addConnection("FLOW", row1Struct.class); } @@ -131,19 +131,28 @@ void elementListenerShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { inputsHandler.setInputValue("FLOW", inputRow); // Force GC to get a clean baseline - System.gc(); - try { - Thread.sleep(100); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } + gc(); + final long memoryBefore = usedMemoryMB(); // --- Producer emits all records in onElement() --- chunkProcessor.onElement(inputFactory, outputFactory); + // Drain — same as Studio generated code (inside + outside loop, skipped for NO_OUTPUT) + int drainedCount = 0; + if (outputMode != OutputMode.NO_OUTPUT) { + // sim inside of input loop + while (outputsHandler.hasMoreData()) { + outputsHandler.getValue("FLOW"); + drainedCount++; + } + } + chunkProcessor.flush(outputFactory); + // GC to reclaim input record objects from the loop + gc(); + // Measure memory AFTER production, BEFORE drain — valid for all output modes final long memoryAfterProduction = usedMemoryMB(); final long memoryDelta = memoryAfterProduction - memoryBefore; @@ -156,14 +165,7 @@ void elementListenerShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { System.out.println("Memory delta: " + memoryDelta + " MB"); System.out.println("Threshold: " + MAX_MEMORY_DELTA_MB + " MB"); - // Drain — same as Studio generated code (inside + outside loop, skipped for NO_OUTPUT) - int drainedCount = 0; if (outputMode != OutputMode.NO_OUTPUT) { - // sim inside of input loop - while (outputsHandler.hasMoreData()) { - outputsHandler.getValue("FLOW"); - drainedCount++; - } // sim outside of input loop while (outputsHandler.hasMoreData()) { outputsHandler.getValue("FLOW"); @@ -209,9 +211,9 @@ void afterGroupShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { inputsHandler.addConnection("FLOW", row1Struct.class); final OutputsHandler outputsHandler = new OutputsHandler(jsonb, servicesMapper); - if(outputMode != OutputMode.NO_OUTPUT) { + if (outputMode != OutputMode.NO_OUTPUT) { outputsHandler.addConnection("MAIN", row1Struct.class); - if(outputMode == OutputMode.TWO_OUTPUT) { + if (outputMode == OutputMode.TWO_OUTPUT) { outputsHandler.addConnection("REJECT", row1Struct.class); } } @@ -225,15 +227,9 @@ void afterGroupShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { int rejectCount = 0; // Force GC to get a clean baseline - System.gc(); - try { - Thread.sleep(100); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } + gc(); final long memoryBefore = usedMemoryMB(); - // Feed 5 inputs (triggers @AfterGroup) for (int i = 0; i < RECORD_COUNT; i++) { final Record inputRecord = builderFactory.newRecordBuilder() .withString("id", "input-" + i) @@ -262,12 +258,7 @@ void afterGroupShouldNotBufferAllRecordsInMemory(OutputMode outputMode) { chunkProcessor.flush(outputFactory); // GC to reclaim input record objects from the loop - System.gc(); - try { - Thread.sleep(100); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } + gc(); final long memoryAfterGroup = usedMemoryMB(); final long memoryDelta = memoryAfterGroup - memoryBefore; @@ -461,4 +452,14 @@ public void readData(final ObjectInputStream objectInputStream) { throw new UnsupportedOperationException("#readData()"); } } + + private void gc() { + // GC to reclaim input record objects from the loop + System.gc(); + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } } From 0409fa9dcbe38a7e1f0d34c211015f1307643b26 Mon Sep 17 00:00:00 2001 From: Thierry Boileau Date: Tue, 21 Jul 2026 17:49:48 +0200 Subject: [PATCH 12/13] test tbo --- .../api/processor/MultiOutputIterator.java | 14 +++++ .../component/api/processor/TaggedOutput.java | 52 +++++++++++++++++++ .../runtime/output/OutputFactory.java | 17 ++++++ .../runtime/output/ProcessorImpl.java | 11 +++- .../runtime/visitor/ModelVisitor.java | 8 +-- .../component/runtime/di/OutputsHandler.java | 19 ++++++- 6 files changed, 115 insertions(+), 6 deletions(-) create mode 100644 component-api/src/main/java/org/talend/sdk/component/api/processor/MultiOutputIterator.java create mode 100644 component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/MultiOutputIterator.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/MultiOutputIterator.java new file mode 100644 index 0000000000000..ee6c8eed435f4 --- /dev/null +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/MultiOutputIterator.java @@ -0,0 +1,14 @@ +package org.talend.sdk.component.api.processor; + +import java.util.Iterator; + +public interface MultiOutputIterator { + /** + * Split mode: sets a single lazy iterator whose elements are tagged with + * the target output connection name via {@link TaggedOutput}. + * The runtime reads one record at a time and routes it to the matching connection. + * + * @param iterator the tagged iterator routing records to their named outputs + */ + void setIterator(Iterator> iterator); +} diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java new file mode 100644 index 0000000000000..e16f2aff8dbcb --- /dev/null +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java @@ -0,0 +1,52 @@ +/** + * Copyright (C) 2006-2026 Talend Inc. - www.talend.com + * + * Licensed 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.talend.sdk.component.api.processor; + +import lombok.Getter; +import lombok.RequiredArgsConstructor; + +/** + * A record tagged with its target output connection name. + * Used with {@link MultiOutputIterator} to route individual records to + * specific output connections from a single streaming iterator. + * + * @param the record type + */ +@Getter +@RequiredArgsConstructor +public class TaggedOutput { + + /** + * The name of the output connection this record should be routed to. + * Use {@code "__default__"} or {@code "FLOW"} for the default output. + */ + private final String outputName; + + /** The record to emit to the named output. */ + private final T rec; + + /** + * Convenience factory method. + * + * @param outputName the target output connection name + * @param rec the record to emit + * @param the record type + * @return a new TaggedOutput + */ + public static TaggedOutput of(final String outputName, final T rec) { + return new TaggedOutput<>(outputName, record); + } +} \ No newline at end of file diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java index b255c29924b25..cadb48ff172f8 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/OutputFactory.java @@ -15,6 +15,7 @@ */ package org.talend.sdk.component.runtime.output; +import org.talend.sdk.component.api.processor.MultiOutputIterator; import org.talend.sdk.component.api.processor.OutputEmitter; import org.talend.sdk.component.api.processor.OutputIterator; @@ -35,4 +36,20 @@ default OutputIterator createIterator(String name) { throw new UnsupportedOperationException( "OutputIterator is only supported in the Studio DI runtime"); } + + /** + * Creates a {@link MultiOutputIterator} that routes records lazily to one or more + * output connections without buffering. + * + *

+ * Supported only in the Studio DI runtime. + * + * @param the record type + * @return a MultiOutputIterator for lazy streaming + * @throws UnsupportedOperationException if the runtime does not support multi-output iterator mode + */ + default MultiOutputIterator createMultiOutputIterator() { + throw new UnsupportedOperationException( + "MultiOutputIterator is only supported in the Studio DI runtime"); + } } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java index 38aab8af13635..fe7009347ed06 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/output/ProcessorImpl.java @@ -55,6 +55,7 @@ import org.talend.sdk.component.api.processor.ElementListener; import org.talend.sdk.component.api.processor.Input; import org.talend.sdk.component.api.processor.LastGroup; +import org.talend.sdk.component.api.processor.MultiOutputIterator; import org.talend.sdk.component.api.processor.Output; import org.talend.sdk.component.api.processor.OutputIterator; import org.talend.sdk.component.api.service.record.RecordBuilderFactory; @@ -154,8 +155,11 @@ private BiFunction buildProcessParamBuilder final String name = parameter.getAnnotation(Output.class).value(); if (OutputIterator.class == parameter.getType()) { return (inputs, outputs) -> outputs.createIterator(name); + } else if (MultiOutputIterator.class == parameter.getType()) { + return (inputs, outputs) -> outputs.createMultiOutputIterator(); + } else { + return (inputs, outputs) -> outputs.create(name); } - return (inputs, outputs) -> outputs.create(name); } final Class parameterType = parameter.getType(); @@ -172,8 +176,11 @@ private Function toOutputParamBuilder(final Parameter par final String name = parameter.getAnnotation(Output.class).value(); if (OutputIterator.class == parameter.getType()) { return outputs.createIterator(name); + } else if (MultiOutputIterator.class == parameter.getType()) { + return outputs.createMultiOutputIterator(); + } else { + return outputs.create(name); } - return outputs.create(name); }; } diff --git a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java index 06d9dadadaee2..91efcaaa8453d 100644 --- a/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java +++ b/component-runtime-impl/src/main/java/org/talend/sdk/component/runtime/visitor/ModelVisitor.java @@ -44,6 +44,7 @@ import org.talend.sdk.component.api.processor.BeforeGroup; import org.talend.sdk.component.api.processor.ElementListener; import org.talend.sdk.component.api.processor.LastGroup; +import org.talend.sdk.component.api.processor.MultiOutputIterator; import org.talend.sdk.component.api.processor.Output; import org.talend.sdk.component.api.processor.OutputEmitter; import org.talend.sdk.component.api.processor.OutputIterator; @@ -197,7 +198,7 @@ private void validateProcessor(final Class input) { final List invalidParams = Stream.of(m.getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { throw new IllegalArgumentException( - "@Output parameter must be of type OutputEmitter or OutputIterator"); + "@Output parameter must be of type OutputEmitter or OutputIterator or MultiOutputIterator"); } }) .filter(p -> !p.isAnnotationPresent(Output.class)) @@ -246,7 +247,7 @@ private void validateProducer(final Class input, final List afterGrou if (!producers.isEmpty() && Stream.of(producers.get(0).getParameters()).peek(p -> { if (p.isAnnotationPresent(Output.class) && !validOutputParam(p)) { throw new IllegalArgumentException( - "@Output parameter must be of type OutputEmitter or OutputIterator"); + "@Output parameter must be of type OutputEmitter or OutputIterator or MultiOutputIterator"); } }).filter(p -> !p.isAnnotationPresent(Output.class)).count() < 1) { throw new IllegalArgumentException(input + " doesn't have the input parameter on its producer method"); @@ -257,7 +258,8 @@ private boolean validOutputParam(final Parameter p) { if (!(p.getParameterizedType() instanceof ParameterizedType pt)) { return false; } - return OutputEmitter.class == pt.getRawType() || OutputIterator.class == pt.getRawType(); + final Type rawType = pt.getRawType(); + return OutputEmitter.class == rawType || OutputIterator.class == rawType || MultiOutputIterator.class == rawType; } private Stream> getPartitionMapperMethods(final boolean infinite) { diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java index ed5e949f25d53..ee44a3c4e5f27 100644 --- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java +++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java @@ -20,8 +20,10 @@ import javax.json.bind.Jsonb; +import org.talend.sdk.component.api.processor.MultiOutputIterator; import org.talend.sdk.component.api.processor.OutputEmitter; import org.talend.sdk.component.api.processor.OutputIterator; +import org.talend.sdk.component.api.processor.TaggedOutput; import org.talend.sdk.component.api.record.Record; import org.talend.sdk.component.api.record.Schema; import org.talend.sdk.component.runtime.output.OutputFactory; @@ -51,7 +53,7 @@ public OutputEmitter create(final String name) { @Override public OutputIterator createIterator(final String name) { final BaseIOHandler.IO ref = connections.get(getActualName(name)); - return iterator -> { + return iterator -> { // wrap the given iterator to provide converted values via the IO. if (ref == null) { return; } @@ -69,6 +71,21 @@ public Object next() { }); }; } + + @Override + public MultiOutputIterator createMultiOutputIterator() { + return taggedOutPutIterator -> { + if (taggedOutPutIterator.hasNext()) { + final TaggedOutput next = taggedOutPutIterator.next(); + final BaseIOHandler.IO ref = connections.get(next.getOutputName()); + final T value = next.getRec(); + + if (ref != null && value != null) { + ref.add(convert(value, ref)); + } + } + }; + } }; } From fe16973c46d8ea1a83f17d1facc6e3be7031e0a3 Mon Sep 17 00:00:00 2001 From: Thierry Boileau Date: Tue, 21 Jul 2026 18:37:16 +0200 Subject: [PATCH 13/13] test tbo --- .../component/api/processor/TaggedOutput.java | 2 +- .../component/runtime/di/OutputsHandler.java | 39 +++++++++++++------ 2 files changed, 28 insertions(+), 13 deletions(-) diff --git a/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java b/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java index e16f2aff8dbcb..c29da1bcfe65d 100644 --- a/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java +++ b/component-api/src/main/java/org/talend/sdk/component/api/processor/TaggedOutput.java @@ -47,6 +47,6 @@ public class TaggedOutput { * @return a new TaggedOutput */ public static TaggedOutput of(final String outputName, final T rec) { - return new TaggedOutput<>(outputName, record); + return new TaggedOutput<>(outputName, rec); } } \ No newline at end of file diff --git a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java index ee44a3c4e5f27..2fec543a1829d 100644 --- a/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java +++ b/component-studio/component-runtime-di/src/main/java/org/talend/sdk/component/runtime/di/OutputsHandler.java @@ -74,21 +74,37 @@ public Object next() { @Override public MultiOutputIterator createMultiOutputIterator() { - return taggedOutPutIterator -> { - if (taggedOutPutIterator.hasNext()) { - final TaggedOutput next = taggedOutPutIterator.next(); - final BaseIOHandler.IO ref = connections.get(next.getOutputName()); - final T value = next.getRec(); - - if (ref != null && value != null) { - ref.add(convert(value, ref)); - } - } - }; + return topi -> setTaggedOutPutIterator(topi); } }; } + private void setTaggedOutPutIterator(final Iterator> taggedOutPutIterator) { + this.taggedOutPutIterator = taggedOutPutIterator; + } + + private Iterator> taggedOutPutIterator; + + @Override + public boolean hasMoreData() { + if (taggedOutPutIterator != null) { + if (taggedOutPutIterator.hasNext()) { + final TaggedOutput next = taggedOutPutIterator.next(); + final IO ref = connections.get(getActualName(next.getOutputName())); + final Object value = next.getRec(); + + if (ref != null && value != null) { + ref.add(convert(value, ref)); + } + return true; + } else { + return false; + } + } else { + return super.hasMoreData(); + } + } + /** * Guess schema special use-case for processor Studio mock. * Same as asOutputFactory but stores the record'schema or schema as the pojo class isn't available. @@ -123,5 +139,4 @@ private Object convert(final Object value, final BaseIOHandler.IO ref) { return jsonb.fromJson(jsonb.toJson(value), ref.getType()); } } - }