Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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) {
Expand Down Expand Up @@ -113,7 +113,7 @@ public MessageProducer getInitializedProducer(SubmissionMode submissionMode) {
}

@Override
public java.util.Queue<PublishResult> getPublishedResultsQueue() {
public BlockingQueue<PublishResult> getPublishedResultsQueue() {
return publishedResultsQueue;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -41,11 +41,11 @@
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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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");

Expand Down Expand Up @@ -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<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);
Expand All @@ -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<Solace.Record> records) {
Expand All @@ -148,17 +161,16 @@ private void publishBatch(List<Solace.Record> 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);
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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");

Expand Down Expand Up @@ -132,7 +134,14 @@ public SessionService solaceSessionServiceWithProducer() {
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;
Expand All @@ -143,9 +152,9 @@ public void publishResults(BeamContextWrapper context) {
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) {
Expand All @@ -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()) {
Expand Down Expand Up @@ -218,6 +230,37 @@ public void publishResults(BeamContextWrapper context) {
}
}

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()) {
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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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")
Expand All @@ -84,6 +88,13 @@ public UnboundedStreamingSolaceWriter(
publishLatencyMetrics);
}

@StartBundle
@Override
public void startBundle() {
super.startBundle();
messageIdsToAck.clear();
}

@ProcessElement
public void processElement(
@Element KV<Integer, Solace.Record> element,
Expand All @@ -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()
Expand Down Expand Up @@ -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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -51,7 +51,7 @@ public MessageProducer getInitializedProducer(SubmissionMode mode) {
}

@Override
public Queue<PublishResult> getPublishedResultsQueue() {
public BlockingQueue<PublishResult> getPublishedResultsQueue() {
throw new UnsupportedOperationException(exceptionMessage);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
}
Loading
Loading