Skip to content

Commit 7db008e

Browse files
committed
Support ACKNOWLEDGEMENT_MODE for JmsIO
* CheckpointMark behavior in alignment with different ACKNOWLEDGEMENT_MODE * Ref count active checkpoint for quicker onClose that releases session * Fix hanging checkpoint when no incoming data in direct runner. This allows us to do an exact assert * Optimize long running unit test usign a short retry * Re-enable AMQP integration test after stuck unack messages resolved
1 parent 401844f commit 7db008e

6 files changed

Lines changed: 570 additions & 159 deletions

File tree

runners/direct-java/src/main/java/org/apache/beam/runners/direct/UnboundedReadEvaluatorFactory.java

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -173,16 +173,17 @@ public void processElement(
173173
} else {
174174
Instant watermark = reader.getWatermark();
175175
if (watermark.isBefore(BoundedWindow.TIMESTAMP_MAX_VALUE)) {
176-
// If the reader had no elements available, but the shard is not done, reuse it later
177-
// Might be better to finalize old checkpoint.
176+
// If the reader had no elements available, but the shard is not done, reuse it later.
177+
// Finalize old checkpoint now.
178+
final CheckpointMarkT checkpoint = shard.getCheckpoint();
179+
if (checkpoint != null) {
180+
checkpoint.finalizeCheckpoint();
181+
}
178182
resultBuilder.addUnprocessedElements(
179183
Collections.<WindowedValue<?>>singleton(
180184
WindowedValues.timestampedValueInGlobalWindow(
181185
UnboundedSourceShard.of(
182-
shard.getSource(),
183-
shard.getDeduplicator(),
184-
reader,
185-
shard.getCheckpoint()),
186+
shard.getSource(), shard.getDeduplicator(), reader, null),
186187
watermark)));
187188
} else {
188189
// End of input. Close the reader after finalizing old checkpoint.

sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsCheckpointMark.java

Lines changed: 84 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -19,13 +19,17 @@
1919

2020
import java.io.IOException;
2121
import java.io.Serializable;
22+
import java.util.ArrayList;
23+
import java.util.List;
2224
import java.util.Objects;
25+
import java.util.concurrent.atomic.AtomicInteger;
2326
import java.util.concurrent.locks.ReentrantReadWriteLock;
2427
import javax.jms.JMSException;
2528
import javax.jms.Message;
2629
import javax.jms.MessageConsumer;
2730
import javax.jms.Session;
2831
import org.apache.beam.sdk.io.UnboundedSource;
32+
import org.apache.beam.sdk.io.jms.JmsIO.AcknowledgeMode;
2933
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
3034
import org.checkerframework.checker.nullness.qual.Nullable;
3135
import org.joda.time.Instant;
@@ -41,29 +45,32 @@ class JmsCheckpointMark implements UnboundedSource.CheckpointMark, Serializable
4145
private static final Logger LOG = LoggerFactory.getLogger(JmsCheckpointMark.class);
4246

4347
private Instant oldestMessageTimestamp;
44-
private transient @Nullable Message lastMessage;
48+
private transient @Nullable List<Message> messages;
4549
private transient @Nullable MessageConsumer consumer;
4650
private transient @Nullable Session session;
51+
private transient @Nullable AtomicInteger activeCheckpoints;
4752

4853
private JmsCheckpointMark(
4954
Instant oldestMessageTimestamp,
50-
@Nullable Message lastMessage,
55+
@Nullable List<Message> messages,
5156
@Nullable MessageConsumer consumer,
52-
@Nullable Session session) {
57+
@Nullable Session session,
58+
@Nullable AtomicInteger activeCheckpoints) {
5359
this.oldestMessageTimestamp = oldestMessageTimestamp;
54-
this.lastMessage = lastMessage;
60+
this.messages = messages;
5561
this.consumer = consumer;
5662
this.session = session;
63+
this.activeCheckpoints = activeCheckpoints;
5764
}
5865

5966
/** Acknowledge all outstanding message. */
6067
@Override
6168
public void finalizeCheckpoint() {
6269
try {
63-
// Jms spec will implicitly acknowledge _all_ messaged already received by the same
64-
// session if one message in this session is being acknowledged.
65-
if (lastMessage != null) {
66-
lastMessage.acknowledge();
70+
if (messages != null) {
71+
for (Message message : messages) {
72+
message.acknowledge();
73+
}
6774
}
6875
} catch (JMSException e) {
6976
// The effect of this is message not get acknowledged and thus will be redelivered. It is
@@ -93,14 +100,37 @@ public void finalizeCheckpoint() {
93100
LOG.info("Error closing JMS session. It may have already been closed.");
94101
}
95102
}
103+
104+
if (activeCheckpoints != null) {
105+
activeCheckpoints.decrementAndGet();
106+
}
107+
}
108+
109+
@VisibleForTesting
110+
@Nullable
111+
List<Message> getMessages() {
112+
return messages;
113+
}
114+
115+
@VisibleForTesting
116+
@Nullable
117+
Session getSession() {
118+
return session;
119+
}
120+
121+
@VisibleForTesting
122+
@Nullable
123+
MessageConsumer getConsumer() {
124+
return consumer;
96125
}
97126

98127
// set an empty list to messages when deserialize
99128
private void readObject(java.io.ObjectInputStream stream)
100129
throws IOException, ClassNotFoundException {
101130
stream.defaultReadObject();
102-
lastMessage = null;
131+
messages = null;
103132
session = null;
133+
consumer = null;
104134
}
105135

106136
@Override
@@ -120,24 +150,27 @@ public int hashCode() {
120150
return Objects.hash(oldestMessageTimestamp);
121151
}
122152

123-
static Preparer newPreparer() {
124-
return new Preparer();
153+
static Preparer newPreparer(AcknowledgeMode acknowledgeMode) {
154+
return new Preparer(acknowledgeMode);
125155
}
126156

127157
/**
128158
* A class preparing the immutable checkpoint. It is mutable so that new messages can be added.
129159
*/
130160
static class Preparer {
131161
private Instant oldestMessageTimestamp = Instant.now();
132-
private transient @Nullable Message lastMessage = null;
162+
private transient List<Message> messages = new ArrayList<>();
163+
private final AcknowledgeMode acknowledgeMode;
133164

134165
@VisibleForTesting transient boolean discarded = false;
135166

136167
@VisibleForTesting final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
137168

138-
private Preparer() {}
169+
private Preparer(AcknowledgeMode acknowledgeMode) {
170+
this.acknowledgeMode = acknowledgeMode;
171+
}
139172

140-
void add(Message message) throws Exception {
173+
void add(Message message) throws JMSException {
141174
lock.writeLock().lock();
142175
try {
143176
if (discarded) {
@@ -149,7 +182,15 @@ void add(Message message) throws Exception {
149182
if (currentMessageTimestamp.isBefore(oldestMessageTimestamp)) {
150183
oldestMessageTimestamp = currentMessageTimestamp;
151184
}
152-
lastMessage = message;
185+
if (acknowledgeMode == AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE) {
186+
messages.add(message);
187+
} else {
188+
if (messages.isEmpty()) {
189+
messages.add(message);
190+
} else {
191+
messages.set(0, message);
192+
}
193+
}
153194
} finally {
154195
lock.writeLock().unlock();
155196
}
@@ -167,6 +208,7 @@ Instant getOldestMessageTimestamp() {
167208
void discard() {
168209
lock.writeLock().lock();
169210
try {
211+
messages.clear();
170212
this.discarded = true;
171213
} finally {
172214
lock.writeLock().unlock();
@@ -175,20 +217,39 @@ void discard() {
175217

176218
/**
177219
* Create a new checkpoint mark based on the current preparer. This will reset the messages held
178-
* by the preparer, and the owner of the preparer is responsible to create a new Jms session
179-
* after this call.
220+
* by the preparer. If AcknowledgeMode is CLIENT_ACKNOWLEDGE, the owner of the preparer is
221+
* responsible to create a new Jms session after this call.
180222
*/
181-
JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Session session) {
223+
JmsCheckpointMark newCheckpoint(
224+
@Nullable MessageConsumer consumer,
225+
@Nullable Session session,
226+
@Nullable AcknowledgeMode acknowledgeMode,
227+
@Nullable AtomicInteger activeCheckpoints) {
182228
JmsCheckpointMark checkpointMark;
183229
lock.writeLock().lock();
184230
try {
185231
if (discarded) {
186-
lastMessage = null;
232+
messages.clear();
187233
checkpointMark = this.emptyCheckpoint();
188234
} else {
235+
List<Message> messagesCopy = null;
236+
MessageConsumer consumerToPass = null;
237+
Session sessionToPass = null;
238+
if (!messages.isEmpty()) {
239+
messagesCopy = new ArrayList<>(messages);
240+
}
241+
if (acknowledgeMode == AcknowledgeMode.CLIENT_ACKNOWLEDGE) {
242+
consumerToPass = consumer;
243+
sessionToPass = session;
244+
}
189245
checkpointMark =
190-
new JmsCheckpointMark(oldestMessageTimestamp, lastMessage, consumer, session);
191-
lastMessage = null;
246+
new JmsCheckpointMark(
247+
oldestMessageTimestamp,
248+
messagesCopy,
249+
consumerToPass,
250+
sessionToPass,
251+
activeCheckpoints);
252+
messages.clear();
192253
oldestMessageTimestamp = Instant.now();
193254
}
194255
} finally {
@@ -198,11 +259,11 @@ JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Se
198259
}
199260

200261
JmsCheckpointMark emptyCheckpoint() {
201-
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null);
262+
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null, null);
202263
}
203264

204265
boolean isEmpty() {
205-
return lastMessage == null;
266+
return messages.isEmpty();
206267
}
207268
}
208269
}

0 commit comments

Comments
 (0)