GH-112: Add KinesisShardEndedEvent to KclMDChA

Fixes https://github.com/spring-projects/spring-integration-aws/issues/112
This commit is contained in:
abilan
2023-03-31 13:43:59 -04:00
parent 75cd725401
commit 10b94e386d

View File

@@ -56,6 +56,8 @@ import software.amazon.kinesis.processor.StreamTracker;
import software.amazon.kinesis.retrieval.KinesisClientRecord;
import software.amazon.kinesis.retrieval.RetrievalConfig;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.ApplicationEventPublisherAware;
import org.springframework.core.AttributeAccessor;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.log.LogMessage;
@@ -63,6 +65,7 @@ import org.springframework.core.serializer.support.DeserializingConverter;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.core.task.TaskExecutor;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.aws.event.KinesisShardEndedEvent;
import org.springframework.integration.aws.support.AwsHeaders;
import org.springframework.integration.endpoint.MessageProducerSupport;
import org.springframework.integration.mapping.InboundMessageMapper;
@@ -86,7 +89,8 @@ import org.springframework.util.Assert;
*/
@ManagedResource
@IntegrationManagedResource
public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
public class KclMessageDrivenChannelAdapter extends MessageProducerSupport
implements ApplicationEventPublisherAware {
private static final ThreadLocal<AttributeAccessor> attributesHolder = new ThreadLocal<>();
@@ -127,6 +131,8 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
private boolean bindSourceRecord;
private ApplicationEventPublisher applicationEventPublisher;
private volatile Scheduler scheduler;
public KclMessageDrivenChannelAdapter(String... streams) {
@@ -153,6 +159,11 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
this.dynamoDBClient = dynamoDBClient;
}
@Override
public void setApplicationEventPublisher(ApplicationEventPublisher applicationEventPublisher) {
this.applicationEventPublisher = applicationEventPublisher;
}
public void setExecutor(TaskExecutor executor) {
Assert.notNull(executor, "'executor' must not be null.");
this.executor = executor;
@@ -407,6 +418,11 @@ public class KclMessageDrivenChannelAdapter extends MessageProducerSupport {
catch (ShutdownException | InvalidStateException ex) {
logger.error(ex, "Exception while checkpointing at requested shutdown. Giving up");
}
if (KclMessageDrivenChannelAdapter.this.applicationEventPublisher != null) {
KclMessageDrivenChannelAdapter.this.applicationEventPublisher.publishEvent(
new KinesisShardEndedEvent(KclMessageDrivenChannelAdapter.this, this.shardId));
}
}
@Override