Skip to content

Commit

Permalink
Add template for kclv3 properties
Browse files Browse the repository at this point in the history
  • Loading branch information
ethkatnic committed Oct 31, 2024
1 parent 82868e7 commit a144dfa
Show file tree
Hide file tree
Showing 4 changed files with 155 additions and 14 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,10 @@ public void setMetricsEnabledDimensions(String[] dimensions) {
@Delegate(types = PollingConfigBean.PollingConfigBeanDelegate.class)
private final PollingConfigBean pollingConfig = new PollingConfigBean();

// TODO: temp
@Delegate(types = WorkerUtilizationAwareAssignmentConfigBean.WorkerUtilizationAwareAssignmentConfigBeanDelegate.class)
private final WorkerUtilizationAwareAssignmentConfigBean workerUtilizationAwareAssignmentConfig = new WorkerUtilizationAwareAssignmentConfigBean();

private boolean validateSequenceNumberBeforeCheckpointing;

private long shutdownGraceMillis;
Expand Down Expand Up @@ -370,6 +374,11 @@ private void handleRetrievalConfig(RetrievalConfig retrievalConfig, ConfigsBuild
retrievalMode.builder(this).build(configsBuilder.kinesisClient(), this));
}

// TODO: temp
private void handleLeaseManagementConfig(LeaseManagementConfig leaseManagementConfig){
leaseManagementConfig.workerUtilizationAwareAssignmentConfig(this.workerUtilizationAwareAssignmentConfig.create());
}

private Object adjustKinesisHttpConfiguration(Object builderObj) {
if (builderObj instanceof KinesisAsyncClientBuilder) {
KinesisAsyncClientBuilder builder = (KinesisAsyncClientBuilder) builderObj;
Expand Down Expand Up @@ -449,6 +458,12 @@ ResolvedConfiguration resolvedConfiguration(ShardRecordProcessorFactory shardRec
retrievalConfig);

handleRetrievalConfig(retrievalConfig, configsBuilder);
// TODO: temp
// handleCoodinatorConfig(coordinatorConfig, configsBuilder);
// LeaseManagementConfig.WorkerUtilizationAwareAssignmentConfig workerUtilizationAwareAssignmentConfig =
// this.workerUtilizationAwareAssignmentConfig.create();
// workerUtilizationAwareAssignmentConfig.workerMetricsTableConfig(this.workerMetricsTableConfig.create());
handleLeaseManagementConfig(leaseManagementConfig);

resolveFields(configObjects, null, new HashSet<>(Arrays.asList(ConfigsBuilder.class, PollingConfig.class)));

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
/*
* Copyright 2019 Amazon.com, Inc. or its affiliates.
* Licensed under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package software.amazon.kinesis.multilang.config;

import java.time.Duration;

import lombok.Getter;
import lombok.Setter;
//import software.amazon.awssdk.services.kinesis.KinesisAsyncClient;
import software.amazon.kinesis.leases.LeaseManagementConfig.WorkerUtilizationAwareAssignmentConfig;

@Getter
@Setter
public class WorkerUtilizationAwareAssignmentConfigBean {

interface WorkerUtilizationAwareAssignmentConfigBeanDelegate {
long inMemoryWorkerMetricsCaptureFrequencyMillis = Duration.ofSeconds(1L).toMillis();

void setInMemoryWorkerMetricsCaptureFrequencyMillis(long value);

}

@ConfigurationSettable(configurationClass = WorkerUtilizationAwareAssignmentConfig.class, convertToOptional = true)
private long inMemoryWorkerMetricsCaptureFrequencyMillis;

public WorkerUtilizationAwareAssignmentConfig create() {
WorkerUtilizationAwareAssignmentConfig conf = new WorkerUtilizationAwareAssignmentConfig();
conf.inMemoryWorkerMetricsCaptureFrequencyMillis(this.inMemoryWorkerMetricsCaptureFrequencyMillis);
return conf;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.mockito.Mock;
import org.mockito.runners.MockitoJUnitRunner;
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;
import software.amazon.kinesis.leases.LeaseManagementConfig;
import software.amazon.kinesis.processor.ShardRecordProcessorFactory;
import software.amazon.kinesis.retrieval.fanout.FanOutConfig;
import software.amazon.kinesis.retrieval.polling.PollingConfig;
Expand Down Expand Up @@ -112,16 +113,29 @@ public void testSetLeaseTableDeletionProtectionEnabledToTrue() {
}

@Test
public void testSetLeaseTablePitrEnabledToTrue() {
public void testSetInMemoryWorkerMetricsCaptureFrequencyMillis() {
final long testVal = 100;
MultiLangDaemonConfiguration configuration = baseConfiguration();
configuration.setLeaseTablePitrEnabled(true);
configuration.setInMemoryWorkerMetricsCaptureFrequencyMillis(testVal);

MultiLangDaemonConfiguration.ResolvedConfiguration resolvedConfiguration =
configuration.resolvedConfiguration(shardRecordProcessorFactory);

assertTrue(resolvedConfiguration.leaseManagementConfig.leaseTablePitrEnabled());
LeaseManagementConfig leaseManagementConfig = resolvedConfiguration.leaseManagementConfig;
LeaseManagementConfig.WorkerUtilizationAwareAssignmentConfig config = leaseManagementConfig.workerUtilizationAwareAssignmentConfig();
assertEquals(config.inMemoryWorkerMetricsCaptureFrequencyMillis(), testVal);
}

// @Test
// public void testSetLeaseTablePitrEnabledToTrue() {
// MultiLangDaemonConfiguration configuration = baseConfiguration();
// configuration.setLeaseTablePitrEnabled(true);
//
// MultiLangDaemonConfiguration.ResolvedConfiguration resolvedConfiguration =
// configuration.resolvedConfiguration(shardRecordProcessorFactory);
//
// assertTrue(resolvedConfiguration.leaseManagementConfig.leaseTablePitrEnabled());
// }

@Test
public void testSetLeaseTableDeletionProtectionEnabledToFalse() {
MultiLangDaemonConfiguration configuration = baseConfiguration();
Expand All @@ -133,16 +147,16 @@ public void testSetLeaseTableDeletionProtectionEnabledToFalse() {
assertFalse(resolvedConfiguration.leaseManagementConfig.leaseTableDeletionProtectionEnabled());
}

@Test
public void testSetLeaseTablePitrEnabledToFalse() {
MultiLangDaemonConfiguration configuration = baseConfiguration();
configuration.setLeaseTablePitrEnabled(false);

MultiLangDaemonConfiguration.ResolvedConfiguration resolvedConfiguration =
configuration.resolvedConfiguration(shardRecordProcessorFactory);

assertFalse(resolvedConfiguration.leaseManagementConfig.leaseTablePitrEnabled());
}
// @Test
// public void testSetLeaseTablePitrEnabledToFalse() {
// MultiLangDaemonConfiguration configuration = baseConfiguration();
// configuration.setLeaseTablePitrEnabled(false);
//
// MultiLangDaemonConfiguration.ResolvedConfiguration resolvedConfiguration =
// configuration.resolvedConfiguration(shardRecordProcessorFactory);
//
// assertFalse(resolvedConfiguration.leaseManagementConfig.leaseTablePitrEnabled());
// }

@Test
public void testDefaultRetrievalConfig() {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
/*
* Copyright 2019 Amazon.com, Inc. or its affiliates.
* Licensed under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package software.amazon.kinesis.multilang.config;

import java.util.Optional;

import org.apache.commons.beanutils.BeanUtilsBean;
import org.apache.commons.beanutils.ConvertUtilsBean;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.runners.MockitoJUnitRunner;
import software.amazon.awssdk.services.kinesis.KinesisAsyncClient;
import software.amazon.kinesis.retrieval.polling.PollingConfig;

import static org.hamcrest.CoreMatchers.equalTo;
import static org.junit.Assert.assertThat;

@RunWith(MockitoJUnitRunner.class)
public class WorkerUtilizationAwareAssignmentConfigBeanTest {

@Mock
private KinesisAsyncClient kinesisAsyncClient;

@Test
public void testAllPropertiesTransit() {
PollingConfigBean pollingConfigBean = new PollingConfigBean();
pollingConfigBean.setIdleTimeBetweenReadsInMillis(1000);
pollingConfigBean.setMaxGetRecordsThreadPool(20);
pollingConfigBean.setMaxRecords(5000);
pollingConfigBean.setRetryGetRecordsInSeconds(30);

ConvertUtilsBean convertUtilsBean = new ConvertUtilsBean();
BeanUtilsBean utilsBean = new BeanUtilsBean(convertUtilsBean);

MultiLangDaemonConfiguration multiLangDaemonConfiguration =
new MultiLangDaemonConfiguration(utilsBean, convertUtilsBean);
multiLangDaemonConfiguration.setStreamName("test-stream");

PollingConfig pollingConfig = pollingConfigBean.build(kinesisAsyncClient, multiLangDaemonConfiguration);

assertThat(pollingConfig.kinesisClient(), equalTo(kinesisAsyncClient));
assertThat(pollingConfig.streamName(), equalTo(multiLangDaemonConfiguration.getStreamName()));
assertThat(
pollingConfig.idleTimeBetweenReadsInMillis(),
equalTo(pollingConfigBean.getIdleTimeBetweenReadsInMillis()));
assertThat(
pollingConfig.maxGetRecordsThreadPool(),
equalTo(Optional.of(pollingConfigBean.getMaxGetRecordsThreadPool())));
assertThat(pollingConfig.maxRecords(), equalTo(pollingConfigBean.getMaxRecords()));
assertThat(
pollingConfig.retryGetRecordsInSeconds(),
equalTo(Optional.of(pollingConfigBean.getRetryGetRecordsInSeconds())));
}
}

0 comments on commit a144dfa

Please sign in to comment.