Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@
import software.amazon.kinesis.coordinator.Scheduler;
import software.amazon.kinesis.exceptions.KinesisClientLibDependencyException;
import software.amazon.kinesis.exceptions.ThrottlingException;
import software.amazon.kinesis.metrics.MetricsConfig;
import software.amazon.kinesis.metrics.NullMetricsFactory;
import software.amazon.kinesis.processor.ShardRecordProcessorFactory;
import software.amazon.kinesis.retrieval.RetrievalConfig;
import software.amazon.kinesis.retrieval.polling.PollingConfig;
Expand Down Expand Up @@ -198,14 +200,18 @@ public Scheduler createScheduler(final Buffer<Record<Event>> buffer) {
kinesisSourceConfig.getPollingConfig().getIdleTimeBetweenReads().toMillis()));
}

MetricsConfig metricsConfig = kinesisSourceConfig.isKclMetricsEnabled()
? configsBuilder.metricsConfig()
: configsBuilder.metricsConfig().metricsFactory(new NullMetricsFactory());

return new Scheduler(
configsBuilder.checkpointConfig(),
configsBuilder.coordinatorConfig()
.schedulerInitializationBackoffTimeMillis(kinesisSourceConfig.getInitializationBackoffTime().toMillis())
.maxInitializationAttempts(kinesisSourceConfig.getMaxInitializationAttempts()),
configsBuilder.leaseManagementConfig().billingMode(BillingMode.PAY_PER_REQUEST),
configsBuilder.lifecycleConfig(),
configsBuilder.metricsConfig(),
metricsConfig,
configsBuilder.processorConfig(),
retrievalConfig
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,10 @@ public Duration getShardAcknowledgmentTimeout() {
@Getter
@JsonProperty("initialization_backoff_time")
private Duration initializationBackoffTime = DEFAULT_INITIALIZATION_BACKOFF_TIME;

@Getter
@JsonProperty("kcl_metrics_enabled")
private boolean kclMetricsEnabled = true;
}


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,9 @@
import software.amazon.awssdk.services.kinesis.model.StreamDescriptionSummary;
import software.amazon.kinesis.common.InitialPositionInStream;
import software.amazon.kinesis.coordinator.Scheduler;
import software.amazon.kinesis.metrics.CloudWatchMetricsFactory;
import software.amazon.kinesis.metrics.MetricsLevel;
import software.amazon.kinesis.metrics.NullMetricsFactory;
import software.amazon.kinesis.retrieval.polling.PollingConfig;

import java.time.Duration;
Expand All @@ -58,6 +60,7 @@

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
Expand Down Expand Up @@ -257,6 +260,30 @@ void testCreateScheduler() {
verify(workerIdentifierGenerator, times(1)).generate();
}

@Test
void testCreateSchedulerUsesNullMetricsFactoryWhenMetricsDisabled() {
when(kinesisSourceConfig.isKclMetricsEnabled()).thenReturn(false);
KinesisService kinesisService = new KinesisService(kinesisSourceConfig, kinesisClientFactory, pluginMetrics, pluginFactory,
pipelineDescription, acknowledgementSetManager, kinesisLeaseConfigSupplier, workerIdentifierGenerator);
Scheduler schedulerObjectUnderTest = kinesisService.createScheduler(buffer);

assertNotNull(schedulerObjectUnderTest);
assertNotNull(schedulerObjectUnderTest.metricsConfig());
assertInstanceOf(NullMetricsFactory.class, schedulerObjectUnderTest.metricsConfig().metricsFactory());
}

@Test
void testCreateSchedulerUsesCloudWatchMetricsFactoryWhenMetricsEnabled() {
when(kinesisSourceConfig.isKclMetricsEnabled()).thenReturn(true);
KinesisService kinesisService = new KinesisService(kinesisSourceConfig, kinesisClientFactory, pluginMetrics, pluginFactory,
pipelineDescription, acknowledgementSetManager, kinesisLeaseConfigSupplier, workerIdentifierGenerator);
Scheduler schedulerObjectUnderTest = kinesisService.createScheduler(buffer);

assertNotNull(schedulerObjectUnderTest);
assertNotNull(schedulerObjectUnderTest.metricsConfig());
assertInstanceOf(CloudWatchMetricsFactory.class, schedulerObjectUnderTest.metricsConfig().metricsFactory());
}

@Test
void testCreateSchedulerWithPollingStrategy() {
when(kinesisSourceConfig.getConsumerStrategy()).thenReturn(ConsumerStrategy.POLLING);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ public class KinesisSourceConfigTest {
private static final String PIPELINE_CONFIG_CHECKPOINT_ENABLED = "pipeline_with_checkpoint_enabled.yaml";
private static final String PIPELINE_CONFIG_STREAM_ARN_ENABLED = "pipeline_with_stream_arn_config.yaml";
private static final String PIPELINE_CONFIG_STREAM_ARN_CONSUMER_ARN_ENABLED = "pipeline_with_stream_arn_consumer_arn_config.yaml";
private static final String PIPELINE_CONFIG_WITH_METRICS_ENABLED = "pipeline_with_metrics_enabled.yaml";
private static final String PIPELINE_CONFIG_WITH_INITIAL_POSITION_AT_TIMESTAMP = "pipeline_with_initial_position_at_timestamp_config.yaml";
private static final Duration MINIMAL_CHECKPOINT_INTERVAL = Duration.ofMillis(2 * 60 * 1000); // 2 minute

Expand Down Expand Up @@ -84,6 +85,7 @@ void testSourceConfig() {
assertEquals(KinesisSourceConfig.DEFAULT_MAX_INITIALIZATION_ATTEMPTS, kinesisSourceConfig.getMaxInitializationAttempts());
assertEquals(KinesisSourceConfig.DEFAULT_INITIALIZATION_BACKOFF_TIME, kinesisSourceConfig.getInitializationBackoffTime());
assertTrue(kinesisSourceConfig.isAcknowledgments());
assertTrue(kinesisSourceConfig.isKclMetricsEnabled());
assertEquals(KinesisSourceConfig.DEFAULT_SHARD_ACKNOWLEDGEMENT_TIMEOUT, kinesisSourceConfig.getShardAcknowledgmentTimeout());
assertThat(kinesisSourceConfig.getAwsAuthenticationConfig(), notNullValue());
assertEquals(kinesisSourceConfig.getAwsAuthenticationConfig().getAwsRegion(), Region.US_EAST_1);
Expand Down Expand Up @@ -117,6 +119,7 @@ void testSourceConfigWithStreamCodec() {
assertEquals(KinesisSourceConfig.DEFAULT_MAX_INITIALIZATION_ATTEMPTS, kinesisSourceConfig.getMaxInitializationAttempts());
assertEquals(KinesisSourceConfig.DEFAULT_INITIALIZATION_BACKOFF_TIME, kinesisSourceConfig.getInitializationBackoffTime());
assertFalse(kinesisSourceConfig.isAcknowledgments());
assertTrue(kinesisSourceConfig.isKclMetricsEnabled());
assertEquals(KinesisSourceConfig.DEFAULT_SHARD_ACKNOWLEDGEMENT_TIMEOUT, kinesisSourceConfig.getShardAcknowledgmentTimeout());
assertThat(kinesisSourceConfig.getAwsAuthenticationConfig(), notNullValue());
assertEquals(kinesisSourceConfig.getAwsAuthenticationConfig().getAwsRegion(), Region.US_EAST_1);
Expand Down Expand Up @@ -239,6 +242,26 @@ void testSourceConfigWithStreamArnConsumerArn() {
}

@Test
@Tag(PIPELINE_CONFIG_WITH_METRICS_ENABLED)
void testSourceConfigWithMetricsEnabled() {

assertThat(kinesisSourceConfig, notNullValue());
assertTrue(kinesisSourceConfig.isKclMetricsEnabled());
assertEquals(KinesisSourceConfig.DEFAULT_NUMBER_OF_RECORDS_TO_ACCUMULATE, kinesisSourceConfig.getNumberOfRecordsToAccumulate());
assertEquals(KinesisSourceConfig.DEFAULT_TIME_OUT_IN_MILLIS, kinesisSourceConfig.getBufferTimeout());
assertFalse(kinesisSourceConfig.isAcknowledgments());
assertThat(kinesisSourceConfig.getAwsAuthenticationConfig(), notNullValue());
assertEquals(kinesisSourceConfig.getAwsAuthenticationConfig().getAwsRegion(), Region.US_EAST_1);
assertEquals(kinesisSourceConfig.getAwsAuthenticationConfig().getAwsStsRoleArn(), "arn:aws:iam::123456789012:role/OSI-PipelineRole");

List<KinesisStreamConfig> streamConfigs = kinesisSourceConfig.getStreams();
assertNotNull(kinesisSourceConfig.getCodec());
assertEquals(streamConfigs.size(), 3);

for (KinesisStreamConfig kinesisStreamConfig: streamConfigs) {
assertTrue(kinesisStreamConfig.getName().contains("stream"));
}
}
@Tag(PIPELINE_CONFIG_WITH_INITIAL_POSITION_AT_TIMESTAMP)
void testSourceConfigWithInitialPositionAtTimestamp() {

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
# Copyright OpenSearch Contributors
# SPDX-License-Identifier: Apache-2.0
#
# The OpenSearch Contributors require contributions made to
# this file be licensed under the Apache-2.0 license or a
# compatible open source license.

source:
kinesis:
streams:
- stream_name: "stream-1"
- stream_name: "stream-2"
- stream_name: "stream-3"
codec:
ndjson:
aws:
sts_role_arn: "arn:aws:iam::123456789012:role/OSI-PipelineRole"
region: "us-east-1"
kcl_metrics_enabled: true

Loading