diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/CheckpointMode.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/CheckpointMode.java index c40ebf7..b3d4deb 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/CheckpointMode.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/CheckpointMode.java @@ -1,5 +1,5 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2017-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -20,6 +20,7 @@ package org.springframework.integration.aws.inbound.kinesis; * The listener mode, record or batch. * * @author Artem Bilan + * @author Hervé Fortin * * @since 1.1 */ @@ -39,6 +40,12 @@ public enum CheckpointMode { /** * Checkpoint on demand via provided to the message {@link Checkpointer} callback. */ - manual + manual, + + /** + * Checkpoint at fixed time intervals. + * @since 2.2.0 + */ + periodic } diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java index ba81fd4..44ddbc1 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KclMessageDrivenChannelAdapter.java @@ -99,7 +99,9 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { private int consumerBackoff; - private long checkpointsInterval = 60_000L; + private long checkpointsInterval = 5_000L; + + private CheckpointMode checkpointMode = CheckpointMode.batch; public KclMessageDrivenChannelAdapter(String streams) { this(streams, AmazonKinesisClientBuilder.defaultClient(), @@ -163,13 +165,17 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { /** * Sets the interval between 2 checkpoints. - * * @param checkpointsInterval interval between 2 checkpoints (in milliseconds) */ public void setCheckpointsInterval(long checkpointsInterval) { this.checkpointsInterval = checkpointsInterval; } + public void setCheckpointMode(CheckpointMode checkpointMode) { + Assert.notNull(checkpointMode, "'checkpointMode' must not be null"); + this.checkpointMode = checkpointMode; + } + @Override protected void onInit() { super.onInit(); @@ -277,7 +283,7 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { } for (Record record : records) { try { - processSingleRecord(record); + processSingleRecord(record, checkpointer); } catch (Throwable t) { logger.warn("Caught throwable while processing record " + record, t); @@ -291,18 +297,37 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { } } } + + // checkpoint if needed + if (CheckpointMode.batch.equals(KclMessageDrivenChannelAdapter.this.checkpointMode)) { + checkpoint(checkpointer); + } + else if (CheckpointMode.periodic.equals(KclMessageDrivenChannelAdapter.this.checkpointMode) && + System.currentTimeMillis() > nextCheckpointTimeInMillis) { + + checkpoint(checkpointer); + this.nextCheckpointTimeInMillis = System.currentTimeMillis() + checkpointsInterval; + } + } /** * Process a single record. * @param record The record to be processed. + * @param checkpointer the checkpointer to use if the checkpointMode is record */ - private void processSingleRecord(Record record) { + private void processSingleRecord(Record record, IRecordProcessorCheckpointer checkpointer) { // Convert AWS Record in Spring Message. - performSend(prepareMessageForRecord(record), record); + performSend(prepareMessageForRecord(record, checkpointer), record); + + // checkpoint if needed + if (CheckpointMode.record.equals(KclMessageDrivenChannelAdapter.this.checkpointMode)) { + checkpoint(checkpointer); + } } - private AbstractIntegrationMessageBuilder prepareMessageForRecord(Record record) { + private AbstractIntegrationMessageBuilder prepareMessageForRecord(Record record, + IRecordProcessorCheckpointer checkpointer) { Object payload = record.getData().array(); Message messageToUse = null; @@ -331,6 +356,10 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport { messageBuilder.copyHeadersIfAbsent(messageToUse.getHeaders()); } + if (CheckpointMode.manual.equals(KclMessageDrivenChannelAdapter.this.checkpointMode)) { + messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, checkpointer); + } + return messageBuilder; } diff --git a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java index 39e0663..08e7ff5 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/kinesis/KinesisMessageDrivenChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2018 the original author or authors. + * Copyright 2017-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -83,6 +83,7 @@ import com.amazonaws.services.kinesis.model.StreamStatus; * * @author Artem Bilan * @author Krzysztof Witkowski + * @author Hervé Fortin * * @since 1.1 */ @@ -138,6 +139,8 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i private CheckpointMode checkpointMode = CheckpointMode.batch; + private long checkpointsInterval = 5_000L; + private int recordsLimit = 10000; private int idleBetweenPolls = 1000; @@ -229,6 +232,15 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i this.checkpointMode = checkpointMode; } + /** + * Sets the interval between 2 checkpoints. Only used when checkpointMode is periodic. + * @param checkpointsInterval interval between 2 checkpoints (in milliseconds) + * @since 2.2.0 + */ + public void setCheckpointsInterval(long checkpointsInterval) { + this.checkpointsInterval = checkpointsInterval; + } + /** * The maximum record to poll per on get-records request. * Not greater then {@code 10000}. @@ -764,6 +776,8 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i private final ShardCheckpointer checkpointer; + private long nextCheckpointTimeInMillis; + private final Runnable processTask = processTask(); private final String key; @@ -1019,6 +1033,11 @@ public class KinesisMessageDrivenChannelAdapter extends MessageProducerSupport i if (CheckpointMode.batch.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) { this.checkpointer.checkpoint(); } + else if (CheckpointMode.periodic.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode) && + System.currentTimeMillis() > nextCheckpointTimeInMillis) { + this.checkpointer.checkpoint(); + this.nextCheckpointTimeInMillis = System.currentTimeMillis() + checkpointsInterval; + } } private AbstractIntegrationMessageBuilder prepareMessageForRecord(Record record) {