This is an automated email from the ASF dual-hosted git repository.

stankiewicz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new ab0ee639755 [SolaceIO] refactor SolaceIO's writers finishBundle 
(#39612)
ab0ee639755 is described below

commit ab0ee639755b7908534882d8bf73b768b380e3dd
Author: RadosÅ‚aw Stankiewicz <[email protected]>
AuthorDate: Thu Sep 17 17:46:52 2026 +0200

    [SolaceIO] refactor SolaceIO's writers finishBundle (#39612)
    
    * refactor SolaceIO's writers to block in @FinishBundle until all published 
persistent messages have received either an acknowledgment (ACK) or a negative 
acknowledgment (NACK/error) from the Solace broker, up to a timeout.
    
    * Transition to blocking queue
    
    * remove timer, introduce state variable to still force checkpoint.
    
    * throw exception if there are non ack-ed messages to avoid data loss
---
 .../sdk/io/solace/broker/JcsmpSessionService.java  |  8 +-
 .../sdk/io/solace/broker/PublishResultHandler.java |  6 +-
 .../beam/sdk/io/solace/broker/SessionService.java  |  4 +-
 .../solace/write/UnboundedBatchedSolaceWriter.java | 76 ++++++++++---------
 .../sdk/io/solace/write/UnboundedSolaceWriter.java | 52 ++++++++++++-
 .../write/UnboundedStreamingSolaceWriter.java      | 21 +++++-
 .../sdk/io/solace/MockEmptySessionService.java     |  4 +-
 .../apache/beam/sdk/io/solace/MockProducer.java    | 54 ++++++++++++++
 .../beam/sdk/io/solace/MockSessionService.java     |  8 +-
 .../sdk/io/solace/MockSessionServiceFactory.java   | 20 ++++-
 .../beam/sdk/io/solace/SolaceIOWriteTest.java      | 87 +++++++++++++++++++++-
 11 files changed, 283 insertions(+), 57 deletions(-)

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 818368a92b9..13db1e606ba 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.Queue;
 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<PublishResult> publishedResultsQueue =
-      new ConcurrentLinkedQueue<>();
+  private final BlockingQueue<PublishResult> publishedResultsQueue = new 
LinkedBlockingQueue<>();
   private final RetryCallableManager retryCallableManager = 
RetryCallableManager.create();
 
   public static JcsmpSessionService create(JCSMPProperties jcsmpProperties, 
@Nullable Queue queue) {
@@ -113,7 +113,7 @@ public abstract class JcsmpSessionService extends 
SessionService {
   }
 
   @Override
-  public java.util.Queue<PublishResult> getPublishedResultsQueue() {
+  public BlockingQueue<PublishResult> 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 1153bfcb7a1..b492ae887e8 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 @@ package org.apache.beam.sdk.io.solace.broker;
 
 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 @@ import org.slf4j.LoggerFactory;
 public final class PublishResultHandler implements 
JCSMPStreamingPublishCorrelatingEventHandler {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(PublishResultHandler.class);
-  private final Queue<PublishResult> publishResultsQueue;
+  private final BlockingQueue<PublishResult> publishResultsQueue;
   private final Counter batchesRejectedByBroker =
       Metrics.counter(UnboundedSolaceWriter.class, "batches_rejected");
 
-  public PublishResultHandler(Queue<PublishResult> publishResultsQueue) {
+  public PublishResultHandler(BlockingQueue<PublishResult> 
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 13aa2808abf..4cb473437d6 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 @@ package org.apache.beam.sdk.io.solace.broker;
 
 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<PublishResult> getPublishedResultsQueue();
+  public abstract BlockingQueue<PublishResult> 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 dd4f81eeb08..81a6af94328 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 @@ package org.apache.beam.sdk.io.solace.write;
 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;
@@ -28,13 +30,11 @@ import org.apache.beam.sdk.io.solace.data.Solace;
 import org.apache.beam.sdk.io.solace.data.Solace.Record;
 import org.apache.beam.sdk.metrics.Counter;
 import org.apache.beam.sdk.metrics.Metrics;
-import org.apache.beam.sdk.state.TimeDomain;
-import org.apache.beam.sdk.state.Timer;
-import org.apache.beam.sdk.state.TimerSpec;
-import org.apache.beam.sdk.state.TimerSpecs;
+import org.apache.beam.sdk.state.StateSpec;
+import org.apache.beam.sdk.state.StateSpecs;
+import org.apache.beam.sdk.state.ValueState;
 import org.apache.beam.sdk.transforms.SerializableFunction;
 import org.apache.beam.sdk.values.KV;
-import org.joda.time.Duration;
 import org.joda.time.Instant;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -64,18 +64,16 @@ 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");
 
   private final Counter batchesRejectedByBroker =
       Metrics.counter(UnboundedSolaceWriter.class, "batches_rejected");
 
-  // State variables are never explicitly "used"
+  // We use a state variable to force a shuffling and ensure the cardinality 
of the processing
   @SuppressWarnings("UnusedVariable")
-  @TimerId("bundle_flusher")
-  private final TimerSpec bundleFlusherTimerSpec = 
TimerSpecs.timer(TimeDomain.PROCESSING_TIME);
+  @StateId("current_key")
+  private final StateSpec<ValueState<Integer>> currentKeySpec = 
StateSpecs.value();
 
   public UnboundedBatchedSolaceWriter(
       SerializableFunction<Record, Destination> destinationFn,
@@ -97,29 +95,39 @@ public final class UnboundedBatchedSolaceWriter extends 
UnboundedSolaceWriter {
   @ProcessElement
   public void processElement(
       @Element KV<Integer, Solace.Record> element,
-      @TimerId("bundle_flusher") Timer bundleFlusherTimer,
-      @Timestamp Instant timestamp) {
+      @Timestamp Instant timestamp,
+      @AlwaysFetched @StateId("current_key") ValueState<Integer> 
currentKeyState) {
 
     setCurrentBundleTimestamp(timestamp);
-
+    Integer currentKey = currentKeyState.read();
+    Integer elementKey = element.getKey();
     Solace.Record record = element.getValue();
 
+    if (currentKey == null || !currentKey.equals(elementKey)) {
+      currentKeyState.write(elementKey);
+    }
+
     if (record == null) {
       LOG.error(
           "SolaceIO.Write: Found null record with key {}. Ignoring record.", 
element.getKey());
     } else {
       addToCurrentBundle(record);
-      // Extend timer for bundle flushing
-      bundleFlusherTimer
-          .offset(Duration.standardSeconds(ACKS_FLUSHING_INTERVAL_SECS))
-          .setRelative();
     }
   }
 
   @FinishBundle
   public void finishBundle(FinishBundleContext context) throws IOException {
-    // Take messages in groups of 50 (if there are enough messages)
     List<Solace.Record> currentBundle = getCurrentBundle();
+    Set<String> 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<Solace.Record> batch = currentBundle.subList(i, toIndex);
@@ -130,12 +138,11 @@ public final class UnboundedBatchedSolaceWriter extends 
UnboundedSolaceWriter {
     }
     getCurrentBundle().clear();
 
-    publishResults(BeamContextWrapper.of(context));
-  }
-
-  @OnTimer("bundle_flusher")
-  public void flushBundle(OnTimerContext context) throws IOException {
-    publishResults(BeamContextWrapper.of(context));
+    if (getDeliveryMode() == DeliveryMode.PERSISTENT && messageIdsToAck != 
null) {
+      waitForAcks(BeamContextWrapper.of(context), messageIdsToAck);
+    } else {
+      publishResults(BeamContextWrapper.of(context), null);
+    }
   }
 
   private void publishBatch(List<Solace.Record> records) {
@@ -148,17 +155,16 @@ public final class UnboundedBatchedSolaceWriter extends 
UnboundedSolaceWriter {
       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 1c98113c241..76e4692b858 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.io.IOException;
 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 abstract class UnboundedSolaceWriter
         currentBundleProducerIndex, sessionServiceFactory, 
writerTransformUuid);
   }
 
-  public void publishResults(BeamContextWrapper context) {
+  public void publishResults(BeamContextWrapper context, @Nullable Set<String> 
messageIdsToAck) {
+    publishResults(context, null, messageIdsToAck);
+  }
+
+  public void publishResults(
+      BeamContextWrapper context,
+      @Nullable PublishResult firstResult,
+      @Nullable Set<String> messageIdsToAck) {
     long sumPublish = 0;
     long countPublish = 0;
     long minPublish = Long.MAX_VALUE;
@@ -143,9 +152,9 @@ public abstract class UnboundedSolaceWriter
     long minFailed = Long.MAX_VALUE;
     long maxFailed = 0;
 
-    Queue<PublishResult> publishResultsQueue =
+    BlockingQueue<PublishResult> 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 abstract class UnboundedSolaceWriter
     }
 
     while (result != null) {
+      if (messageIdsToAck != null) {
+        messageIdsToAck.remove(result.getMessageId());
+      }
       Long latency = result.getLatencyNanos();
 
       if (latency == null && shouldPublishLatencyMetrics()) {
@@ -218,6 +230,38 @@ public abstract class UnboundedSolaceWriter
     }
   }
 
+  public void waitForAcks(BeamContextWrapper context, Set<String> 
messageIdsToAck) {
+    BlockingQueue<PublishResult> 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()) {
+      String errorMessage =
+          String.format(
+              "SolaceIO.Write: Timed out waiting for ACKs of %d messages. 
Outstanding message IDs: %s",
+              messageIdsToAck.size(), messageIdsToAck);
+      throw new RuntimeException(errorMessage);
+    }
+  }
+
   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 6d6d0b27e2b..0db0ee9047a 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 @@ package org.apache.beam.sdk.io.solace.write;
 
 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<String> 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 final class UnboundedStreamingSolaceWriter extends 
UnboundedSolaceWriter
         publishLatencyMetrics);
   }
 
+  @StartBundle
+  @Override
+  public void startBundle() {
+    super.startBundle();
+    messageIdsToAck.clear();
+  }
+
   @ProcessElement
   public void processElement(
       @Element KV<Integer, Solace.Record> element,
@@ -105,6 +116,10 @@ public final class UnboundedStreamingSolaceWriter extends 
UnboundedSolaceWriter
       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 final class UnboundedStreamingSolaceWriter extends 
UnboundedSolaceWriter
 
   @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 f6bb6741954..cf014060c2c 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 @@ package org.apache.beam.sdk.io.solace;
 
 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 abstract class MockEmptySessionService extends 
SessionService {
   }
 
   @Override
-  public Queue<PublishResult> getPublishedResultsQueue() {
+  public BlockingQueue<PublishResult> 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 27131035957..a1712633535 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 abstract class MockProducer implements 
MessageProducer {
       }
     }
   }
+
+  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 e888c62c852..e9b780ed0ba 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.google.auto.value.AutoValue;
 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<PublishResultHandler, MockProducer> 
mockProducerFn();
 
-  private final Queue<PublishResult> publishedResultsReceiver = new 
ConcurrentLinkedQueue<>();
+  private final BlockingQueue<PublishResult> publishedResultsReceiver = new 
LinkedBlockingQueue<>();
 
   public static Builder builder() {
     return new AutoValue_MockSessionService.Builder()
@@ -94,7 +94,7 @@ public abstract class MockSessionService extends 
SessionService {
   }
 
   @Override
-  public Queue<PublishResult> getPublishedResultsQueue() {
+  public BlockingQueue<PublishResult> 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 9c17ca60420..5844cd2a741 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 @@ package org.apache.beam.sdk.io.solace;
 
 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 abstract class MockSessionServiceFactory extends 
SessionServiceFactory {
             .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 abstract class MockSessionServiceFactory extends 
SessionServiceFactory {
   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 3cdc392fa1f..c67458b2c98 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
@@ -115,8 +115,21 @@ public class SolaceIOWriteTest {
       WriterType writerType,
       Pipeline p,
       ErrorHandler<BadRecord, ?> errorHandler) {
+    return getWriteTransform(
+        mode, writerType, p, errorHandler, 
SessionServiceType.WITH_SUCCEEDING_PRODUCER);
+  }
+
+  private SolaceOutput getWriteTransform(
+      SubmissionMode mode,
+      WriterType writerType,
+      Pipeline p,
+      ErrorHandler<BadRecord, ?> errorHandler,
+      SessionServiceType sessionServiceType) {
     SessionServiceFactory fakeSessionServiceFactory =
-        MockSessionServiceFactory.builder().mode(mode).build();
+        MockSessionServiceFactory.builder()
+            .mode(mode)
+            .sessionServiceType(sessionServiceType)
+            .build();
 
     PCollection<Record> records = getRecords(p);
     return records.apply(
@@ -288,4 +301,76 @@ public class SolaceIOWriteTest {
         .isEqualTo((long) payloads.size());
     pipeline.run();
   }
+
+  @Test
+  public void testWriteLatencyStreamingWithDelayedAck() throws Exception {
+    SubmissionMode mode = SubmissionMode.LOWER_LATENCY;
+    WriterType writerType = WriterType.STREAMING;
+
+    ErrorHandler<BadRecord, PCollection<Long>> errorHandler =
+        pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform());
+    SolaceOutput output =
+        getWriteTransform(
+            mode, writerType, pipeline, errorHandler, 
SessionServiceType.WITH_DELAYED_PRODUCER);
+    PCollection<String> 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<BadRecord, PCollection<Long>> errorHandler =
+        pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform());
+    SolaceOutput output =
+        getWriteTransform(
+            mode, writerType, pipeline, errorHandler, 
SessionServiceType.WITH_DELAYED_PRODUCER);
+    PCollection<String> 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<BadRecord, PCollection<Long>> errorHandler =
+        pipeline.registerBadRecordErrorHandler(new ErrorSinkTransform());
+
+    SessionServiceFactory fakeSessionServiceFactory =
+        MockSessionServiceFactory.builder()
+            .mode(mode)
+            .sessionServiceType(SessionServiceType.WITH_EXCEPTION_PRODUCER)
+            .build();
+
+    PCollection<Record> 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<String> ids = getIdsPCollection(output);
+
+    PAssert.that(ids).empty();
+    errorHandler.close();
+    PAssert.thatSingleton(Objects.requireNonNull(errorHandler.getOutput()))
+        .isEqualTo((long) payloads.size());
+    pipeline.run();
+  }
 }

Reply via email to