diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/JcsmpSessionService.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/JcsmpSessionService.java index 818368a92b9f..13db1e606ba7 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/JcsmpSessionService.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/JcsmpSessionService.java @@ -32,8 +32,9 @@ import com.solacesystems.jcsmp.XMLMessageProducer; import java.io.IOException; import java.util.Objects; +import java.util.concurrent.BlockingQueue; import java.util.concurrent.Callable; -import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.LinkedBlockingQueue; import javax.annotation.Nullable; import org.apache.beam.sdk.io.solace.RetryCallableManager; import org.apache.beam.sdk.io.solace.SolaceIO.SubmissionMode; @@ -57,8 +58,7 @@ public abstract class JcsmpSessionService extends SessionService { @Nullable private transient JCSMPSession jcsmpSession; @Nullable private transient MessageReceiver messageReceiver; @Nullable private transient MessageProducer messageProducer; - private final java.util.Queue publishedResultsQueue = - new ConcurrentLinkedQueue<>(); + private final BlockingQueue publishedResultsQueue = new LinkedBlockingQueue<>(); private final RetryCallableManager retryCallableManager = RetryCallableManager.create(); public static JcsmpSessionService create(JCSMPProperties jcsmpProperties, @Nullable Queue queue) { @@ -113,7 +113,7 @@ public MessageProducer getInitializedProducer(SubmissionMode submissionMode) { } @Override - public java.util.Queue getPublishedResultsQueue() { + public BlockingQueue getPublishedResultsQueue() { return publishedResultsQueue; } diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/PublishResultHandler.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/PublishResultHandler.java index 1153bfcb7a1c..b492ae887e85 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/PublishResultHandler.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/PublishResultHandler.java @@ -19,7 +19,7 @@ import com.solacesystems.jcsmp.JCSMPException; import com.solacesystems.jcsmp.JCSMPStreamingPublishCorrelatingEventHandler; -import java.util.Queue; +import java.util.concurrent.BlockingQueue; import org.apache.beam.sdk.io.solace.data.Solace; import org.apache.beam.sdk.io.solace.data.Solace.PublishResult; import org.apache.beam.sdk.io.solace.write.UnboundedSolaceWriter; @@ -41,11 +41,11 @@ public final class PublishResultHandler implements JCSMPStreamingPublishCorrelatingEventHandler { private static final Logger LOG = LoggerFactory.getLogger(PublishResultHandler.class); - private final Queue publishResultsQueue; + private final BlockingQueue publishResultsQueue; private final Counter batchesRejectedByBroker = Metrics.counter(UnboundedSolaceWriter.class, "batches_rejected"); - public PublishResultHandler(Queue publishResultsQueue) { + public PublishResultHandler(BlockingQueue publishResultsQueue) { this.publishResultsQueue = publishResultsQueue; } diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/SessionService.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/SessionService.java index 13aa2808abf0..4cb473437d6f 100644 --- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/SessionService.java +++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/broker/SessionService.java @@ -19,7 +19,7 @@ import com.solacesystems.jcsmp.JCSMPProperties; import java.io.Serializable; -import java.util.Queue; +import java.util.concurrent.BlockingQueue; import org.apache.beam.sdk.io.solace.SolaceIO; import org.apache.beam.sdk.io.solace.SolaceIO.SubmissionMode; import org.apache.beam.sdk.io.solace.data.Solace.PublishResult; @@ -138,7 +138,7 @@ public abstract class SessionService implements Serializable { * asynchronously received callbacks from Solace for message publications. The queue * implementation has to be thread-safe for production use-cases. */ - public abstract Queue getPublishedResultsQueue(); + public abstract BlockingQueue getPublishedResultsQueue(); /** * Override this method and provide your specific properties, including all those related to 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..4293315ad050 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 @@ -29,8 +29,9 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; -import java.util.Queue; +import java.util.Set; import java.util.UUID; +import java.util.concurrent.BlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import org.apache.beam.sdk.annotations.Internal; @@ -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,14 @@ public SessionService solaceSessionServiceWithProducer() { currentBundleProducerIndex, sessionServiceFactory, writerTransformUuid); } - public void publishResults(BeamContextWrapper context) { + public void publishResults(BeamContextWrapper context, @Nullable Set messageIdsToAck) { + publishResults(context, null, messageIdsToAck); + } + + public void publishResults( + BeamContextWrapper context, + @Nullable PublishResult firstResult, + @Nullable Set messageIdsToAck) { long sumPublish = 0; long countPublish = 0; long minPublish = Long.MAX_VALUE; @@ -143,9 +152,9 @@ public void publishResults(BeamContextWrapper context) { long minFailed = Long.MAX_VALUE; long maxFailed = 0; - Queue publishResultsQueue = + BlockingQueue publishResultsQueue = solaceSessionServiceWithProducer().getPublishedResultsQueue(); - Solace.PublishResult result = publishResultsQueue.poll(); + PublishResult result = firstResult != null ? firstResult : publishResultsQueue.poll(); if (result != null) { if (getCurrentBundleTimestamp() == null) { @@ -154,6 +163,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 +230,37 @@ public void publishResults(BeamContextWrapper context) { } } + public void waitForAcks(BeamContextWrapper context, Set messageIdsToAck) { + BlockingQueue queue = + solaceSessionServiceWithProducer().getPublishedResultsQueue(); + long timeoutMs = System.currentTimeMillis() + ACKS_FLUSHING_INTERVAL_SECS * 1000; + while (!messageIdsToAck.isEmpty()) { + publishResults(context, messageIdsToAck); + if (messageIdsToAck.isEmpty()) { + break; + } + long remainingTimeMs = timeoutMs - System.currentTimeMillis(); + if (remainingTimeMs <= 0) { + break; + } + try { + PublishResult result = queue.poll(remainingTimeMs, TimeUnit.MILLISECONDS); + if (result != null) { + publishResults(context, result, messageIdsToAck); + } + } 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/MockEmptySessionService.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockEmptySessionService.java index f6bb67419541..cf014060c2ce 100644 --- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockEmptySessionService.java +++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockEmptySessionService.java @@ -19,7 +19,7 @@ import com.google.auto.value.AutoValue; import com.solacesystems.jcsmp.JCSMPProperties; -import java.util.Queue; +import java.util.concurrent.BlockingQueue; import org.apache.beam.sdk.io.solace.SolaceIO.SubmissionMode; import org.apache.beam.sdk.io.solace.broker.MessageProducer; import org.apache.beam.sdk.io.solace.broker.MessageReceiver; @@ -51,7 +51,7 @@ public MessageProducer getInitializedProducer(SubmissionMode mode) { } @Override - public Queue getPublishedResultsQueue() { + public BlockingQueue getPublishedResultsQueue() { throw new UnsupportedOperationException(exceptionMessage); } 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/MockSessionService.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionService.java index e888c62c8522..e9b780ed0ba8 100644 --- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionService.java +++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/MockSessionService.java @@ -21,8 +21,8 @@ import com.solacesystems.jcsmp.BytesXMLMessage; import com.solacesystems.jcsmp.JCSMPProperties; import java.io.IOException; -import java.util.Queue; -import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Function; import org.apache.beam.sdk.io.solace.MockProducer.MockSuccessProducer; @@ -48,7 +48,7 @@ public abstract class MockSessionService extends SessionService { public abstract Function mockProducerFn(); - private final Queue publishedResultsReceiver = new ConcurrentLinkedQueue<>(); + private final BlockingQueue publishedResultsReceiver = new LinkedBlockingQueue<>(); public static Builder builder() { return new AutoValue_MockSessionService.Builder() @@ -94,7 +94,7 @@ public MessageProducer getInitializedProducer(SubmissionMode mode) { } @Override - public Queue getPublishedResultsQueue() { + public BlockingQueue getPublishedResultsQueue() { return publishedResultsReceiver; } 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(); + } }