Skip to content

Commit 66bd2c9

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 66bd2c9

6 files changed

Lines changed: 573 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: 87 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,18 @@ 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+
// Jms spec will implicitly acknowledge _all_ messaged already received by the same
189+
// session if one message in this session is being acknowledged. Only need to ack
190+
// last seen one.
191+
if (messages.isEmpty()) {
192+
messages.add(message);
193+
} else {
194+
messages.set(0, message);
195+
}
196+
}
153197
} finally {
154198
lock.writeLock().unlock();
155199
}
@@ -167,6 +211,7 @@ Instant getOldestMessageTimestamp() {
167211
void discard() {
168212
lock.writeLock().lock();
169213
try {
214+
messages.clear();
170215
this.discarded = true;
171216
} finally {
172217
lock.writeLock().unlock();
@@ -175,20 +220,39 @@ void discard() {
175220

176221
/**
177222
* 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.
223+
* by the preparer. If AcknowledgeMode is CLIENT_ACKNOWLEDGE, the owner of the preparer is
224+
* responsible to create a new Jms session after this call.
180225
*/
181-
JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Session session) {
226+
JmsCheckpointMark newCheckpoint(
227+
@Nullable MessageConsumer consumer,
228+
@Nullable Session session,
229+
@Nullable AcknowledgeMode acknowledgeMode,
230+
@Nullable AtomicInteger activeCheckpoints) {
182231
JmsCheckpointMark checkpointMark;
183232
lock.writeLock().lock();
184233
try {
185234
if (discarded) {
186-
lastMessage = null;
235+
messages.clear();
187236
checkpointMark = this.emptyCheckpoint();
188237
} else {
238+
List<Message> messagesCopy = null;
239+
MessageConsumer consumerToPass = null;
240+
Session sessionToPass = null;
241+
if (!messages.isEmpty()) {
242+
messagesCopy = new ArrayList<>(messages);
243+
}
244+
if (acknowledgeMode == AcknowledgeMode.CLIENT_ACKNOWLEDGE) {
245+
consumerToPass = consumer;
246+
sessionToPass = session;
247+
}
189248
checkpointMark =
190-
new JmsCheckpointMark(oldestMessageTimestamp, lastMessage, consumer, session);
191-
lastMessage = null;
249+
new JmsCheckpointMark(
250+
oldestMessageTimestamp,
251+
messagesCopy,
252+
consumerToPass,
253+
sessionToPass,
254+
activeCheckpoints);
255+
messages.clear();
192256
oldestMessageTimestamp = Instant.now();
193257
}
194258
} finally {
@@ -198,11 +262,11 @@ JmsCheckpointMark newCheckpoint(@Nullable MessageConsumer consumer, @Nullable Se
198262
}
199263

200264
JmsCheckpointMark emptyCheckpoint() {
201-
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null);
265+
return new JmsCheckpointMark(oldestMessageTimestamp, null, null, null, null);
202266
}
203267

204268
boolean isEmpty() {
205-
return lastMessage == null;
269+
return messages.isEmpty();
206270
}
207271
}
208272
}

0 commit comments

Comments
 (0)