Skip to content

Commit aa7f7e8

Browse files
committed
Address comments
1 parent 66bd2c9 commit aa7f7e8

5 files changed

Lines changed: 42 additions & 26 deletions

File tree

CHANGES.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
## New Features / Improvements
7070

7171
* (Python) Removed the `envoy-data-plane` (and transitive `betterproto`) dependency; `EnvoyRateLimiter` now uses a small vendored protobuf definition instead, resolving dependency conflicts for downstream projects ([#37854](https://github.com/apache/beam/issues/37854)).
72+
* (Java) Supported acknowledge mode for JmsIO ([#39253](https://github.com/apache/beam/issues/39253)).
7273
* X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
7374

7475
## Breaking Changes

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -254,6 +254,9 @@ JmsCheckpointMark newCheckpoint(
254254
activeCheckpoints);
255255
messages.clear();
256256
oldestMessageTimestamp = Instant.now();
257+
if (activeCheckpoints != null) {
258+
activeCheckpoints.incrementAndGet();
259+
}
257260
}
258261
} finally {
259262
lock.writeLock().unlock();

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

Lines changed: 26 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import java.util.Optional;
3535
import java.util.UUID;
3636
import java.util.concurrent.ScheduledExecutorService;
37+
import java.util.concurrent.TimeUnit;
3738
import java.util.concurrent.atomic.AtomicInteger;
3839
import java.util.stream.Stream;
3940
import javax.jms.Connection;
@@ -510,7 +511,12 @@ public Read<T> withRequiresDeduping() {
510511
return builder().setRequiresDeduping(true).build();
511512
}
512513

513-
/** Specify the {@link AcknowledgeMode} used for consuming and acknowledging JMS messages. */
514+
/**
515+
* Specify the {@link AcknowledgeMode} used for consuming and acknowledging JMS messages.
516+
*
517+
* <p>To use {@link AcknowledgeMode#INDIVIDUAL_ACKNOWLEDGE}, providers other than ActiveMQ, Qpid
518+
* JMS, ActiveMQ require configuring {@link #withIndividualAcknowledgeModeCode} explicitly.
519+
*/
514520
public Read<T> withAcknowledgeMode(AcknowledgeMode acknowledgeMode) {
515521
checkArgument(acknowledgeMode != null, "acknowledgeMode can not be null");
516522
return builder().setAcknowledgeMode(acknowledgeMode).build();
@@ -521,7 +527,7 @@ public Read<T> withAcknowledgeMode(AcknowledgeMode acknowledgeMode) {
521527
* AcknowledgeMode#INDIVIDUAL_ACKNOWLEDGE}.
522528
*
523529
* <p>Different JMS providers use different proprietary integer constants for individual
524-
* acknowledgment (e.g., ActiveMQ uses 4, Qpid JMS / ActiveMQ Artemis / IBM MQ use 101).
530+
* acknowledgment (e.g., ActiveMQ uses 4, Qpid JMS / ActiveMQ Artemis use 101).
525531
*/
526532
public Read<T> withIndividualAcknowledgeModeCode(int individualAcknowledgeModeCode) {
527533
return builder().setIndividualAcknowledgeModeCode(individualAcknowledgeModeCode).build();
@@ -872,7 +878,6 @@ public CheckpointMark getCheckpointMark() {
872878
sessionTofinalize = null;
873879
}
874880
}
875-
activeCheckpoints.incrementAndGet();
876881
return checkpointMarkPreparer.newCheckpoint(
877882
consumerToClose, sessionTofinalize, mode, activeCheckpoints);
878883
}
@@ -911,25 +916,26 @@ private void doClose() {
911916
} else {
912917
ScheduledExecutorService executorService =
913918
options.as(ExecutorOptions.class).getScheduledExecutorService();
914-
executorService.submit(
915-
() -> {
916-
long startTime = System.currentTimeMillis();
917-
long timeoutMillis = source.spec.getCloseTimeout().getMillis();
918-
while (activeCheckpoints.get() > 0
919-
&& System.currentTimeMillis() - startTime < timeoutMillis) {
920-
try {
921-
Thread.sleep(1_000); // poll in 1 sec interval
922-
} catch (InterruptedException ignored) {
923-
break;
919+
long deadline = System.currentTimeMillis() + source.spec.getCloseTimeout().getMillis();
920+
long pollInterval = 1L;
921+
executorService.schedule(
922+
new Runnable() {
923+
@Override
924+
public void run() {
925+
if (activeCheckpoints.get() == 0 || System.currentTimeMillis() >= deadline) {
926+
LOG.debug(
927+
"Closing connection after checkpoints finalized or timeout: {}",
928+
source.spec.getCloseTimeout());
929+
closeConsumer();
930+
closeSession();
931+
closeConnection();
932+
} else {
933+
executorService.schedule(this, pollInterval, TimeUnit.SECONDS);
924934
}
925935
}
926-
LOG.debug(
927-
"Closing connection after checkpoints finalized or timeout: {}",
928-
source.spec.getCloseTimeout());
929-
closeConsumer();
930-
closeSession();
931-
closeConnection();
932-
});
936+
},
937+
pollInterval,
938+
TimeUnit.SECONDS);
933939
}
934940
} catch (Exception e) {
935941
LOG.warn("Error closing reader", e);

sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOIT.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -293,7 +293,7 @@ private void runPublishingThenReadingAll(JmsIO.AcknowledgeMode acknowledgeMode)
293293

294294
private void cancelIfTimeouted(PipelineResult readResult, PipelineResult.State readState)
295295
throws IOException {
296-
if (readState == null) {
296+
if (readState == null || !readState.isTerminal()) {
297297
readResult.cancel();
298298
}
299299
}

sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOTest.java

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,8 @@
4343
import static org.junit.Assert.assertTrue;
4444
import static org.junit.Assert.fail;
4545
import static org.mockito.ArgumentMatchers.any;
46+
import static org.mockito.ArgumentMatchers.anyLong;
47+
import static org.mockito.ArgumentMatchers.eq;
4648
import static org.mockito.Mockito.mock;
4749
import static org.mockito.Mockito.times;
4850
import static org.mockito.Mockito.verify;
@@ -65,6 +67,7 @@
6567
import java.util.List;
6668
import java.util.Set;
6769
import java.util.concurrent.ScheduledExecutorService;
70+
import java.util.concurrent.TimeUnit;
6871
import java.util.concurrent.atomic.AtomicInteger;
6972
import java.util.function.Function;
7073
import javax.jms.BytesMessage;
@@ -694,7 +697,7 @@ public void testJmsCheckpointMarkIndividualAcknowledgeAllMessages() throws Excep
694697
preparer.add(msg2);
695698
preparer.add(msg3);
696699

697-
AtomicInteger activeCheckpoints = new AtomicInteger(1);
700+
AtomicInteger activeCheckpoints = new AtomicInteger(0);
698701
JmsCheckpointMark mark =
699702
preparer.newCheckpoint(
700703
null, null, JmsIO.AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE, activeCheckpoints);
@@ -722,7 +725,7 @@ public void testJmsCheckpointMarkClientAcknowledgeUnsafeNoSessionRecreation() th
722725
preparer.add(msg1);
723726
preparer.add(msg2);
724727

725-
AtomicInteger activeCheckpoints = new AtomicInteger(1);
728+
AtomicInteger activeCheckpoints = new AtomicInteger(0);
726729
JmsCheckpointMark mark =
727730
preparer.newCheckpoint(
728731
null, null, JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE_UNSAFE, activeCheckpoints);
@@ -960,17 +963,20 @@ public void testCloseWithTimeout() throws IOException, JMSException {
960963
ExecutorOptions options = PipelineOptionsFactory.as(ExecutorOptions.class);
961964
options.setScheduledExecutorService(mockScheduledExecutorService);
962965
ArgumentCaptor<Runnable> runnableArgumentCaptor = ArgumentCaptor.forClass(Runnable.class);
963-
when(mockScheduledExecutorService.submit(runnableArgumentCaptor.capture()))
966+
when(mockScheduledExecutorService.schedule(
967+
runnableArgumentCaptor.capture(), anyLong(), any(TimeUnit.class)))
964968
.thenReturn(null /* unused */);
965969

966970
JmsIO.UnboundedJmsReader reader = source.createReader(options, null);
967971
reader.start();
968972
assertFalse(getDiscardedValue(reader));
969973
reader.checkpointMarkPreparer.add(Mockito.mock(Message.class));
970-
reader.getCheckpointMark();
974+
CheckpointMark mark = reader.getCheckpointMark();
971975
reader.close();
972976
assertTrue(getDiscardedValue(reader));
973-
verify(mockScheduledExecutorService).submit(any(Runnable.class));
977+
verify(mockScheduledExecutorService)
978+
.schedule(any(Runnable.class), eq(1L), eq(TimeUnit.SECONDS));
979+
mark.finalizeCheckpoint();
974980
runnableArgumentCaptor.getValue().run();
975981
assertTrue(getDiscardedValue(reader));
976982
verifyNoMoreInteractions(mockScheduledExecutorService);

0 commit comments

Comments
 (0)