Skip to content

Commit a5736f9

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 a5736f9

6 files changed

Lines changed: 576 additions & 159 deletions

File tree

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

Lines changed: 8 additions & 7 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.
@@ -257,7 +258,7 @@ private CheckpointMarkT finishRead(
257258
final CheckpointMark oldMark = shard.getCheckpoint();
258259
@SuppressWarnings("unchecked")
259260
final CheckpointMarkT mark = (CheckpointMarkT) reader.getCheckpointMark();
260-
if (oldMark != null) {
261+
if (oldMark != null && oldMark != mark) {
261262
oldMark.finalizeCheckpoint();
262263
}
263264

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

Lines changed: 92 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,31 @@ public int hashCode() {
120150
return Objects.hash(oldestMessageTimestamp);
121151
}
122152

123-
static Preparer newPreparer() {
124-
return new Preparer();
153+
protected static Preparer newPreparer() {
154+
return new Preparer(AcknowledgeMode.CLIENT_ACKNOWLEDGE);
155+
}
156+
157+
protected static Preparer newPreparer(@Nullable AcknowledgeMode acknowledgeMode) {
158+
return new Preparer(acknowledgeMode);
125159
}
126160

127161
/**
128162
* A class preparing the immutable checkpoint. It is mutable so that new messages can be added.
129163
*/
130164
static class Preparer {
131165
private Instant oldestMessageTimestamp = Instant.now();
132-
private transient @Nullable Message lastMessage = null;
166+
private transient List<Message> messages = new ArrayList<>();
167+
private final @Nullable AcknowledgeMode acknowledgeMode;
133168

134169
@VisibleForTesting transient boolean discarded = false;
135170

136171
@VisibleForTesting final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
137172

138-
private Preparer() {}
173+
private Preparer(@Nullable AcknowledgeMode acknowledgeMode) {
174+
this.acknowledgeMode = acknowledgeMode;
175+
}
139176

140-
void add(Message message) throws Exception {
177+
void add(Message message) throws JMSException {
141178
lock.writeLock().lock();
142179
try {
143180
if (discarded) {
@@ -149,7 +186,15 @@ void add(Message message) throws Exception {
149186
if (currentMessageTimestamp.isBefore(oldestMessageTimestamp)) {
150187
oldestMessageTimestamp = currentMessageTimestamp;
151188
}
152-
lastMessage = message;
189+
if (acknowledgeMode == AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE) {
190+
messages.add(message);
191+
} else {
192+
if (messages.isEmpty()) {
193+
messages.add(message);
194+
} else {
195+
messages.set(0, message);
196+
}
197+
}
153198
} finally {
154199
lock.writeLock().unlock();
155200
}
@@ -167,28 +212,52 @@ Instant getOldestMessageTimestamp() {
167212
void discard() {
168213
lock.writeLock().lock();
169214
try {
215+
messages.clear();
170216
this.discarded = true;
171217
} finally {
172218
lock.writeLock().unlock();
173219
}
174220
}
175221

222+
JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Session session) {
223+
return newCheckpoint(consumer, session, JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE, null);
224+
}
225+
176226
/**
177227
* 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.
228+
* by the preparer. If AcknowledgeMode is CLIENT_ACKNOWLEDGE, the owner of the preparer is
229+
* responsible to create a new Jms session after this call.
180230
*/
181-
JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Session session) {
231+
JmsCheckpointMark newCheckpoint(
232+
@Nullable MessageConsumer consumer,
233+
@Nullable Session session,
234+
@Nullable AcknowledgeMode acknowledgeMode,
235+
@Nullable AtomicInteger activeCheckpoints) {
182236
JmsCheckpointMark checkpointMark;
183237
lock.writeLock().lock();
184238
try {
185239
if (discarded) {
186-
lastMessage = null;
240+
messages.clear();
187241
checkpointMark = this.emptyCheckpoint();
188242
} else {
243+
List<Message> messagesCopy = null;
244+
MessageConsumer consumerToPass = null;
245+
Session sessionToPass = null;
246+
if (!messages.isEmpty()) {
247+
messagesCopy = new ArrayList<>(messages);
248+
}
249+
if (acknowledgeMode == AcknowledgeMode.CLIENT_ACKNOWLEDGE) {
250+
consumerToPass = consumer;
251+
sessionToPass = session;
252+
}
189253
checkpointMark =
190-
new JmsCheckpointMark(oldestMessageTimestamp, lastMessage, consumer, session);
191-
lastMessage = null;
254+
new JmsCheckpointMark(
255+
oldestMessageTimestamp,
256+
messagesCopy,
257+
consumerToPass,
258+
sessionToPass,
259+
activeCheckpoints);
260+
messages.clear();
192261
oldestMessageTimestamp = Instant.now();
193262
}
194263
} finally {
@@ -198,11 +267,11 @@ JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Se
198267
}
199268

200269
JmsCheckpointMark emptyCheckpoint() {
201-
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null);
270+
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null, null);
202271
}
203272

204273
boolean isEmpty() {
205-
return lastMessage == null;
274+
return messages.isEmpty();
206275
}
207276
}
208277
}

0 commit comments

Comments
 (0)