【问题标题】:How to disable CloudWatch metrics for KPL/KCL with Spring Cloud Stream如何使用 Spring Cloud Stream 禁用 KPL/KCL 的 CloudWatch 指标
【发布时间】:2021-12-07 03:39:51
【问题描述】:

我正在使用 Spring Cloud Stream Binder for Kinesis 并启用了 KPL/KCL。我们希望禁用 Cloudwatch 指标,而不必自己管理 KPL 和 KCL 的配置(完全覆盖 bean)。除了KinesisProducerConfiguration.setMetricsLevel()KinesisClientLibConfiguration.withMetricsLevel(...) 属性之外,我们希望对KinesisProducerConfigurationKinesisClientLibConfiguration 使用相同的bean 定义。

作为参考,这里是在 Spring Cloud Stream Kinesis Binder 中定义 AWS bean 的位置:KinesisBinderConfiguration.java

最有效的方法是什么?

感谢任何帮助!谢谢。

【问题讨论】:

    标签: java spring-boot spring-cloud-stream amazon-kinesis amazon-kinesis-kpl


    【解决方案1】:

    框架不提供任何KinesisClientLibConfiguration。公开这样一个 bean 以及你需要的任何选项是你的项目责任:https://github.com/spring-cloud/spring-cloud-stream-binder-aws-kinesis/blob/main/spring-cloud-stream-binder-kinesis-docs/src/main/asciidoc/overview.adoc#kinesis-consumer-properties

    从版本 2.0.1 开始,可以在应用程序上下文中提供 KinesisClientLibConfiguration 类型的 bean,以完全控制 Kinesis 客户端库配置选项。

    生产者端确实被KinesisBinderConfiguration 中的KinesisProducerConfiguration bean 覆盖:

    @Bean
    @ConditionalOnMissingBean
    @ConditionalOnProperty(name = "spring.cloud.stream.kinesis.binder.kpl-kcl-enabled")
    public KinesisProducerConfiguration kinesisProducerConfiguration() {
        KinesisProducerConfiguration kinesisProducerConfiguration = new KinesisProducerConfiguration();
        kinesisProducerConfiguration.setCredentialsProvider(this.awsCredentialsProvider);
        kinesisProducerConfiguration.setRegion(this.region);
        return kinesisProducerConfiguration;
    }
    

    从这里我看不出有什么大问题,在您自己的配置中声明这样一个 bean 并带有您希望拥有的任何其他属性,包括提到的指标。

    如果这仍然不适合您,您可以将 bean 注入到您自己的 bean 中,并以任何您想要的方式对其进行变异:

    @Bean
    String configurerBean(KinesisProducerConfiguration kinesisProducerConfiguration)  {
       kinesisProducerConfiguration.setMetricsLevel();
       return null;
    }
    

    更新

    消费者部分:

    这是一个基于我们在内部创建的 KCL 的默认配置实例的 bean:

    @Bean
    KinesisClientLibConfiguration kinesisClientLibConfiguration() {
        return new KinesisClientLibConfiguration(this.consumerGroup,
                                this.stream,
                                null,
                                null,
                                this.streamInitialSequence,
                                this.kinesisProxyCredentialsProvider,
                                null,
                                null,
                                KinesisClientLibConfiguration.DEFAULT_FAILOVER_TIME_MILLIS,
                                this.workerId,
                                KinesisClientLibConfiguration.DEFAULT_MAX_RECORDS,
                                this.idleBetweenPolls,
                                false,
                                KinesisClientLibConfiguration.DEFAULT_PARENT_SHARD_POLL_INTERVAL_MILLIS,
                                KinesisClientLibConfiguration.DEFAULT_SHARD_SYNC_INTERVAL_MILLIS,
                                KinesisClientLibConfiguration.DEFAULT_CLEANUP_LEASES_UPON_SHARDS_COMPLETION,
                                new ClientConfiguration(),
                                new ClientConfiguration(),
                                new ClientConfiguration(),
                                this.consumerBackoff,
                                KinesisClientLibConfiguration.DEFAULT_METRICS_BUFFER_TIME_MILLIS,
                                KinesisClientLibConfiguration.DEFAULT_METRICS_MAX_QUEUE_SIZE,
                                KinesisClientLibConfiguration.DEFAULT_VALIDATE_SEQUENCE_NUMBER_BEFORE_CHECKPOINTING,
                                null,
                                KinesisClientLibConfiguration.DEFAULT_SHUTDOWN_GRACE_MILLIS,
                                KinesisClientLibConfiguration.DEFAULT_DDB_BILLING_MODE,
                                new SimpleRecordsFetcherFactory(),
                                DEFAULT_LEASE_CLEANUP_INTERVAL_MILLIS,
                                DEFAULT_COMPLETED_LEASE_CLEANUP_THRESHOLD_MILLIS,
                                DEFAULT_GARBAGE_LEASE_CLEANUP_THRESHOLD_MILLIS);
    }
    

    您在this. 看到的任何内容都必须替换为您环境中的相应值。在这种情况下,KinesisClientLibConfiguration.DEFAULT_METRICS_MAX_QUEUE_SIZE 可能就是您要查找的内容。

    this.consumerGroupthis.stream 必须与您要为其配置使用者的绑定相同。

    【讨论】:

    • 非常感谢您的回复。您介意为 KCL 消费者提供一个示例吗?鉴于许多重载的构造函数已被弃用,我正在努力找出最好的方法
    • 在我的回答中查看更新。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-01-11
    • 2018-01-11
    • 1970-01-01
    • 2018-01-24
    • 1970-01-01
    • 2020-10-26
    • 1970-01-01
    相关资源
    最近更新 更多