diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedBatchedSolaceWriter.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedBatchedSolaceWriter.java index dd4f81eeb082..49e6bd76b858 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedBatchedSolaceWriter.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedBatchedSolaceWriter.java @@ -20,7 +20,9 @@ import com.solacesystems.jcsmp.DeliveryMode; import com.solacesystems.jcsmp.Destination; import java.io.IOException; +import java.util.HashSet; import java.util.List; +import java.util.Set; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.io.solace.SolaceIO.SubmissionMode; import org.apache.beam.sdk.io.solace.broker.SessionServiceFactory; @@ -64,8 +66,6 @@ public final class UnboundedBatchedSolaceWriter extends UnboundedSolaceWriter { private static final Logger LOG = LoggerFactory.getLogger(UnboundedBatchedSolaceWriter.class); - private static final int ACKS_FLUSHING_INTERVAL_SECS = 10; - private final Counter sentToBroker = Metrics.counter(UnboundedBatchedSolaceWriter.class, "msgs_sent_to_broker"); @@ -118,8 +118,17 @@ public void processElement( @FinishBundle public void finishBundle(FinishBundleContext context) throws IOException { - // Take messages in groups of 50 (if there are enough messages) List currentBundle = getCurrentBundle(); + Set messageIdsToAck = null; + + if (getDeliveryMode() == DeliveryMode.PERSISTENT) { + messageIdsToAck = new HashSet<>(); + for (Solace.Record record : currentBundle) { + messageIdsToAck.add(record.getMessageId()); + } + } + + // Take messages in groups of 50 (if there are enough messages) for (int i = 0; i < currentBundle.size(); i += SOLACE_BATCH_LIMIT) { int toIndex = Math.min(i + SOLACE_BATCH_LIMIT, currentBundle.size()); List batch = currentBundle.subList(i, toIndex); @@ -130,12 +139,16 @@ public void finishBundle(FinishBundleContext context) throws IOException { } getCurrentBundle().clear(); - publishResults(BeamContextWrapper.of(context)); + if (getDeliveryMode() == DeliveryMode.PERSISTENT && messageIdsToAck != null) { + waitForAcks(BeamContextWrapper.of(context), messageIdsToAck); + } else { + publishResults(BeamContextWrapper.of(context), null); + } } @OnTimer("bundle_flusher") public void flushBundle(OnTimerContext context) throws IOException { - publishResults(BeamContextWrapper.of(context)); + publishResults(BeamContextWrapper.of(context), null); } private void publishBatch(List records) { @@ -148,17 +161,16 @@ private void publishBatch(List records) { sentToBroker.inc(entriesPublished); } catch (Exception e) { batchesRejectedByBroker.inc(); - Solace.PublishResult errorPublish = - Solace.PublishResult.builder() - .setPublished(false) - .setMessageId(String.format("BATCH_OF_%d_ENTRIES", records.size())) - .setError( - String.format( - "Batch could not be published after several" + " retries. Error: %s", - e.getMessage())) - .setLatencyNanos(System.nanoTime()) - .build(); - solaceSessionServiceWithProducer().getPublishedResultsQueue().add(errorPublish); + for (Solace.Record record : records) { + Solace.PublishResult errorPublish = + Solace.PublishResult.builder() + .setPublished(false) + .setMessageId(record.getMessageId()) + .setError(String.format("Batch could not be published. Error: %s", e.getMessage())) + .setLatencyNanos(System.nanoTime()) + .build(); + solaceSessionServiceWithProducer().getPublishedResultsQueue().add(errorPublish); + } } } } diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedSolaceWriter.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedSolaceWriter.java index 1c98113c2416..c0b47b3089ed 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedSolaceWriter.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedSolaceWriter.java @@ -30,6 +30,7 @@ import java.util.List; import java.util.Optional; import java.util.Queue; +import java.util.Set; import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -68,6 +69,7 @@ public abstract class UnboundedSolaceWriter // This is the batch limit supported by the send multiple JCSMP API method. static final int SOLACE_BATCH_LIMIT = 50; + static final int ACKS_FLUSHING_INTERVAL_SECS = 10; private final Distribution latencyPublish = Metrics.distribution(SolaceIO.Write.class, "latency_publish_ms"); @@ -132,7 +134,7 @@ public SessionService solaceSessionServiceWithProducer() { currentBundleProducerIndex, sessionServiceFactory, writerTransformUuid); } - public void publishResults(BeamContextWrapper context) { + public void publishResults(BeamContextWrapper context, @Nullable Set messageIdsToAck) { long sumPublish = 0; long countPublish = 0; long minPublish = Long.MAX_VALUE; @@ -154,6 +156,9 @@ public void publishResults(BeamContextWrapper context) { } while (result != null) { + if (messageIdsToAck != null) { + messageIdsToAck.remove(result.getMessageId()); + } Long latency = result.getLatencyNanos(); if (latency == null && shouldPublishLatencyMetrics()) { @@ -218,6 +223,27 @@ public void publishResults(BeamContextWrapper context) { } } + public void waitForAcks(BeamContextWrapper context, Set messageIdsToAck) { + long timeoutMs = System.currentTimeMillis() + ACKS_FLUSHING_INTERVAL_SECS * 1000; + while (!messageIdsToAck.isEmpty() && System.currentTimeMillis() < timeoutMs) { + publishResults(context, messageIdsToAck); + if (!messageIdsToAck.isEmpty()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + break; + } + } + } + if (!messageIdsToAck.isEmpty()) { + LOG.warn( + "SolaceIO.Write: Timed out waiting for ACKs of {} messages. Outstanding message IDs: {}", + messageIdsToAck.size(), + messageIdsToAck); + } + } + public BytesXMLMessage createSingleMessage( Solace.Record record, boolean useCorrelationKeyLatency) { JCSMPFactory jcsmpFactory = JCSMPFactory.onlyInstance(); diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedStreamingSolaceWriter.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedStreamingSolaceWriter.java index 6d6d0b27e2bb..0db0ee9047aa 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedStreamingSolaceWriter.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/write/UnboundedStreamingSolaceWriter.java @@ -19,6 +19,8 @@ import com.solacesystems.jcsmp.DeliveryMode; import com.solacesystems.jcsmp.Destination; +import java.util.HashSet; +import java.util.Set; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.io.solace.SolaceIO; import org.apache.beam.sdk.io.solace.broker.SessionServiceFactory; @@ -63,6 +65,8 @@ public final class UnboundedStreamingSolaceWriter extends UnboundedSolaceWriter private final Counter rejectedByBroker = Metrics.counter(UnboundedStreamingSolaceWriter.class, "msgs_rejected_by_broker"); + private final Set messageIdsToAck = new HashSet<>(); + // We use a state variable to force a shuffling and ensure the cardinality of the processing @SuppressWarnings("UnusedVariable") @StateId("current_key") @@ -84,6 +88,13 @@ public UnboundedStreamingSolaceWriter( publishLatencyMetrics); } + @StartBundle + @Override + public void startBundle() { + super.startBundle(); + messageIdsToAck.clear(); + } + @ProcessElement public void processElement( @Element KV element, @@ -105,6 +116,10 @@ public void processElement( return; } + if (getDeliveryMode() == DeliveryMode.PERSISTENT) { + messageIdsToAck.add(record.getMessageId()); + } + // The publish method will retry, let's send a failure message if all the retries fail try { solaceSessionServiceWithProducer() @@ -133,6 +148,10 @@ public void processElement( @FinishBundle public void finishBundle(FinishBundleContext context) { - publishResults(BeamContextWrapper.of(context)); + if (getDeliveryMode() == DeliveryMode.PERSISTENT) { + waitForAcks(BeamContextWrapper.of(context), messageIdsToAck); + } else { + publishResults(BeamContextWrapper.of(context), null); + } } } diff --git a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockProducer.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockProducer.java index 271310359577..a1712633535b 100644 --- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockProducer.java +++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockProducer.java @@ -107,4 +107,58 @@ public void publishSingleMessage( } } } + + public static class MockDelayedProducer extends MockProducer { + private final long delayMs; + + public MockDelayedProducer(PublishResultHandler handler, long delayMs) { + super(handler); + this.delayMs = delayMs; + } + + public MockDelayedProducer(PublishResultHandler handler) { + this(handler, 100); + } + + @Override + public void publishSingleMessage( + Record msg, + Destination topicOrQueue, + boolean useCorrelationKeyLatency, + DeliveryMode deliveryMode) { + new Thread( + () -> { + try { + Thread.sleep(delayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + if (useCorrelationKeyLatency) { + handler.responseReceivedEx( + Solace.PublishResult.builder() + .setPublished(true) + .setMessageId(msg.getMessageId()) + .build()); + } else { + handler.responseReceivedEx(msg.getMessageId()); + } + }) + .start(); + } + } + + public static class MockExceptionProducer extends MockProducer { + public MockExceptionProducer(PublishResultHandler handler) { + super(handler); + } + + @Override + public void publishSingleMessage( + Record msg, + Destination topicOrQueue, + boolean useCorrelationKeyLatency, + DeliveryMode deliveryMode) { + throw new RuntimeException("Simulated synchronous publish failure"); + } + } } diff --git a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionServiceFactory.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionServiceFactory.java index 9c17ca604201..5844cd2a7415 100644 --- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionServiceFactory.java +++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionServiceFactory.java @@ -19,6 +19,8 @@ import com.google.auto.value.AutoValue; import com.solacesystems.jcsmp.BytesXMLMessage; +import org.apache.beam.sdk.io.solace.MockProducer.MockDelayedProducer; +import org.apache.beam.sdk.io.solace.MockProducer.MockExceptionProducer; import org.apache.beam.sdk.io.solace.MockProducer.MockFailedProducer; import org.apache.beam.sdk.io.solace.MockProducer.MockSuccessProducer; import org.apache.beam.sdk.io.solace.SolaceIO.SubmissionMode; @@ -80,6 +82,20 @@ public SessionService create() { .mode(mode()) .mockProducerFn(MockFailedProducer::new) .build(); + case WITH_DELAYED_PRODUCER: + return MockSessionService.builder() + .recordFn(recordFn()) + .minMessagesReceived(minMessagesReceived()) + .mode(mode()) + .mockProducerFn(MockDelayedProducer::new) + .build(); + case WITH_EXCEPTION_PRODUCER: + return MockSessionService.builder() + .recordFn(recordFn()) + .minMessagesReceived(minMessagesReceived()) + .mode(mode()) + .mockProducerFn(MockExceptionProducer::new) + .build(); default: throw new RuntimeException( String.format("Unknown sessionServiceType: %s", sessionServiceType().name())); @@ -89,6 +105,8 @@ public SessionService create() { public enum SessionServiceType { EMPTY, WITH_SUCCEEDING_PRODUCER, - WITH_FAILING_PRODUCER + WITH_FAILING_PRODUCER, + WITH_DELAYED_PRODUCER, + WITH_EXCEPTION_PRODUCER } } diff --git a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/SolaceIOWriteTest.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/SolaceIOWriteTest.java index e92657c3c3d2..55dff02d999b 100644 --- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/SolaceIOWriteTest.java +++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/SolaceIOWriteTest.java @@ -87,8 +87,21 @@ private SolaceOutput getWriteTransform( WriterType writerType, Pipeline p, ErrorHandler errorHandler) { + return getWriteTransform( + mode, writerType, p, errorHandler, SessionServiceType.WITH_SUCCEEDING_PRODUCER); + } + + private SolaceOutput getWriteTransform( + SubmissionMode mode, + WriterType writerType, + Pipeline p, + ErrorHandler errorHandler, + SessionServiceType sessionServiceType) { SessionServiceFactory fakeSessionServiceFactory = - MockSessionServiceFactory.builder().mode(mode).build(); + MockSessionServiceFactory.builder() + .mode(mode) + .sessionServiceType(sessionServiceType) + .build(); PCollection records = getRecords(p); return records.apply( @@ -205,4 +218,76 @@ public void testWriteWithFailedRecords() throws Exception { .isEqualTo((long) payloads.size()); pipeline.run(); } + + @Test + public void testWriteLatencyStreamingWithDelayedAck() throws Exception { + SubmissionMode mode = SubmissionMode.LOWER_LATENCY; + WriterType writerType = WriterType.STREAMING; + + ErrorHandler> errorHandler = + pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform()); + SolaceOutput output = + getWriteTransform( + mode, writerType, pipeline, errorHandler, SessionServiceType.WITH_DELAYED_PRODUCER); + PCollection ids = getIdsPCollection(output); + + PAssert.that(ids).containsInAnyOrder(keys); + errorHandler.close(); + PAssert.that(errorHandler.getOutput()).empty(); + + pipeline.run(); + } + + @Test + public void testWriteLatencyBatchedWithDelayedAck() throws Exception { + SubmissionMode mode = SubmissionMode.LOWER_LATENCY; + WriterType writerType = WriterType.BATCHED; + + ErrorHandler> errorHandler = + pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform()); + SolaceOutput output = + getWriteTransform( + mode, writerType, pipeline, errorHandler, SessionServiceType.WITH_DELAYED_PRODUCER); + PCollection ids = getIdsPCollection(output); + + PAssert.that(ids).containsInAnyOrder(keys); + errorHandler.close(); + PAssert.that(errorHandler.getOutput()).empty(); + + pipeline.run(); + } + + @Test + public void testWriteWithExceptionRecords() throws Exception { + SubmissionMode mode = SubmissionMode.HIGHER_THROUGHPUT; + WriterType writerType = WriterType.BATCHED; + ErrorHandler> errorHandler = + pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform()); + + SessionServiceFactory fakeSessionServiceFactory = + MockSessionServiceFactory.builder() + .mode(mode) + .sessionServiceType(SessionServiceType.WITH_EXCEPTION_PRODUCER) + .build(); + + PCollection records = getRecords(pipeline); + SolaceOutput output = + records.apply( + "Write to Solace", + SolaceIO.write() + .to(Solace.Queue.fromName("queue")) + .withSubmissionMode(mode) + .withWriterType(writerType) + .withDeliveryMode(DeliveryMode.PERSISTENT) + .withSessionServiceFactory(fakeSessionServiceFactory) + .withErrorHandler(errorHandler)); + + PCollection ids = getIdsPCollection(output); + + PAssert.that(ids).empty(); + errorHandler.close(); + PAssert.thatSingleton(Objects.requireNonNull(errorHandler.getOutput())) + .isEqualTo((long) payloads.size()); + pipeline.run(); + } }