【问题标题】:Wrong Kafka topic names for Spring-Cloud-Function app deployed as part of Spring-Cloud-Data-Flow stream作为 Spring-Cloud-Data-Flow 流的一部分部署的 Spring-Cloud-Function 应用程序的错误 Kafka 主题名称
【发布时间】:2020-08-12 12:33:18
【问题描述】:

我有一个简单的 SCDF 流,如下所示:

http --port=12346 | mvmn-transform | file --name=tmp.txt --directory=/tmp

mvmn-transform 是一个简单的自定义转换器,如下所示:

@SpringBootApplication
@EnableBinding(Processor.class)
@EnableConfigurationProperties(ScdfTestTransformerProperties.class)
@Configuration
public class ScdfTestTransformer {
    public static void main(String args[]) {
        SpringApplication.run(ScdfTestTransformer.class, args);
    }

    @Autowired
    protected ScdfTestTransformerProperties config;

    @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT)
    public Object transform(Message<?> message) {
        Object payload = message.getPayload();
        Map<String, Object> result = new HashMap<>();
        Map<String, String> headersStr = new HashMap<>();

        message.getHeaders().forEach((k, v) -> headersStr.put(k, v != null ? v.toString() : null));

        result.put("headers", headersStr);
        result.put("payload", payload);
        result.put("configProp", config.getSomeConfigProp());

        return result;
    }

    // See https://stackoverflow.com/questions/59155689/could-not-decode-json-type-for-key-file-name-in-a-spring-cloud-data-flow-stream
    @Bean("kafkaBinderHeaderMapper")
    public KafkaHeaderMapper kafkaBinderHeaderMapper() {
        BinderHeaderMapper mapper = new BinderHeaderMapper();
        mapper.setEncodeStrings(true);
        return mapper;
    }
}

这很好用。

但我读到 Spring Cloud Function 应该允许我实现这样的应用程序而无需指定绑定和转换器注释,所以我将其更改为:

@SpringBootApplication
// @EnableBinding(Processor.class)
@EnableConfigurationProperties(ScdfTestTransformerProperties.class)
@Configuration
public class ScdfTestTransformer {
    public static void main(String args[]) {
        SpringApplication.run(ScdfTestTransformer.class, args);
    }

    @Autowired
    protected ScdfTestTransformerProperties config;

    // @Transformer(inputChannel = Processor.INPUT, outputChannel = Processor.OUTPUT)
    @Bean
    public Function<Message<?>, Map<String, Object>> transform(
    // Message<?> message
    ) {
        return message -> {
            Object payload = message.getPayload();
            Map<String, Object> result = new HashMap<>();
            Map<String, String> headersStr = new HashMap<>();

            message.getHeaders().forEach((k, v) -> headersStr.put(k, v != null ? v.toString() : null));

            result.put("headers", headersStr);
            result.put("payload", payload);
            result.put("configProp", "Config prop val: " + config.getSomeConfigProp());

            return result;
        };
    }

    // See https://stackoverflow.com/questions/59155689/could-not-decode-json-type-for-key-file-name-in-a-spring-cloud-data-flow-stream
    @Bean("kafkaBinderHeaderMapper")
    public KafkaHeaderMapper kafkaBinderHeaderMapper() {
        BinderHeaderMapper mapper = new BinderHeaderMapper();
        mapper.setEncodeStrings(true);
        return mapper;
    }
}

现在我遇到了一个问题 - Spring-Cloud-Function 显然忽略了 SCDF 源和目标主题名称,而是创建了主题 transform-in-0transform-out-0

SCDF 创建名称类似于&lt;stream-name&gt;.&lt;app-name&gt; 的主题,例如TestStream123.httpTestStream123.mvmn-transform

以前它们被用于转换器 - 应该如此,因为它是 SCDF 流的一部分。但是现在它们被 Spring-Cloud-Function 忽略了,而是创建了 transform-in-0transform-out-0

因此,我的转换器不再接收任何输入,因为它期望它在错误的 Kafka 主题上。并且可能也不会对流产生任何输出,因为它也会输出到错误的 Kafka 主题。

附:以防万一,GitHub上的完整项目代码:https://github.com/mvmn/scdftest-transformer/tree/scfunc

为了在本地运行启动 Kafka、Skipper、SCDF 和 SCDF 控制台,请在 app 文件夹中执行 mvn clean install,然后在控制台中执行 app register --name mvmn-transform-1 --type processor --uri maven://x.mvmn.study.scdf.scdftest:scdftest-transformer:0.1.1-SNAPSHOT --metadata-uri maven://x.mvmn.study.scdf.scdftest:scdftest-transformer:0.1.1-SNAPSHOT。然后你可以使用定义http --port=12346 | mvmn-transform | file --name=tmp.txt --directory=/tmp部署流

【问题讨论】:

    标签: spring-cloud-dataflow spring-cloud-function


    【解决方案1】:

    由于您使用的是编写 Spring Cloud Stream 应用程序的功能模型,因此在部署此应用程序时,您需要在自定义处理器上传递两个属性来恢复 Spring Cloud Data Flow 行为。

    spring.cloud.stream.function.bindings.transform-in-0=input spring.cloud.stream.function.bindings.transform-out-0=output

    你可以试试看,看看会不会有什么不同?

    【讨论】:

    • 非常感谢!我已经把这些放到 application.properties 中,这样部署这个转换器就没有特殊的步骤了——而且效果很好!再次感谢。
    • P.S.假设我们也可以使用 ${spring.cloud.stream.instanceIndex} 而不是 0。但我不确定。
    • P.P.S.从头开始 - 这是函数绑定的索引,而不是应用程序实例 - 在单个函数应用程序中它始终为 0。更多信息在这里:github.com/spring-cloud/spring-cloud-stream/blob/master/docs/… - 以防万一。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-04-06
    • 1970-01-01
    • 2019-01-18
    • 1970-01-01
    • 2022-01-09
    • 2019-07-22
    相关资源
    最近更新 更多