Kafka Connect 与 Spring 框架

Posted

技术标签:

【中文标题】Kafka Connect 与 Spring 框架【英文标题】:Kafka Connect with Spring Framework 【发布时间】:2018-01-16 11:51:22 【问题描述】:

有人知道 Spring Boot 与 Kafka Connect 的集成吗?我认为有一个 spring-kafka 项目可以很好地与 Kafka 客户端集成,但不能连接和流式传输 API。

【问题讨论】:

你想用 kafka connect & spring 做什么? 【参考方案1】:

KConnect 本身是一个不同的 Client,你不能给它添加 spring。通过这样做,您最终会创建自己的客户端,这不能称为 KConnect。如果您打算将 KConnect 插件与 spring 集成,这是一项艰巨的任务,但应该仍然可行,但我不会推荐它,因为插件在初始化和运行时间方面应该很轻。此外,KConnect 并不意味着包含任何业务逻辑,如果没有业务逻辑,那么标准插件应该能够完全满足您的需求。

但是,KStreams 可以与 spring 集成。就像创建一个 bean 对象一样简单。这是一个示例

public class SampleStream 

@Autowired
CustomBean myBean;

private static final Logger LOG = LoggerFactory.getLogger(SampleStream.class);

public SampleStream() 
    KafkaStreams stream = getStream();
    LOG.info("Starting Stream ", stream);

    stream.start();
    Runtime.getRuntime().addShutdownHook(new Thread(stream::close));


private KafkaStreams getStream()

    StreamsBuilder builder = new StreamsBuilder();
    KStream<String, String> requestStream = builder.stream("REQUEST_TOPIC");

    KStream<String, String> responseStream = requestStream.flatMap((key, request) -> 
         myBean.process(request)
        //custom logic
    );

    responseStream.to("RESPONSE_TOPIC");

    return new KafkaStreams(builder.build(), getStreamProperties());



private Properties getStreamProperties() 
    String defaultKeySerdes = Serdes.StringSerde.class.getName();
    String defaultValueSerdes = Serdes.StringSerde.class.getName();
    String defaultExceptionHandler = LogAndContinueExceptionHandler.class.getName();

    Properties properties = new Properties();
    properties.put(StreamsConfig.APPLICATION_ID_CONFIG, this.getClass().toString());
    properties.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka.com:9092");
    properties.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, defaultKeySerdes);
    properties.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, defaultValueSerdes);
    properties.put(GenericMessagingConstants.KAFKA_DESERIALIZER_VALUE_CLASS_CONFIG, String.class.getName());
    properties.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, defaultExceptionHandler);
    ....
    return properties;
    

在应用程序上下文中创建 Bean,如下所示,

<bean id="sampleStream" class="SampleStream"/>

回答这个问题需要一些时间才能解决。

【讨论】:

以上是关于Kafka Connect 与 Spring 框架的主要内容,如果未能解决你的问题,请参考以下文章

如何在没有 Confluent 的情况下使用 Kafka Connect for Cassandra

Github标星5.3K,kafka分区与副本

我的 Python/Java/Spring/Go/Whatever Client Won’t Connect to My Apache Kafka Cluster in Docker/AWS/My B

Kafka Connect HDFS

kafka connect 使用说明

KafkaKIP-285 Connect Rest Extension Plugin kafka 连接 rest 的插件