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