【问题标题】:Configuring what is written to Kafka Topic when using Kafka Connect Cassandra Source配置使用 Kafka Connect Cassandra 源时写入 Kafka 主题的内容
【发布时间】:2017-07-20 02:00:36
【问题描述】:

我正在做一个峰值,我们希望将数据发布到 Cassandra 表中,并将其写入 Kafka 主题。我们正在考虑使用 Kafka Connect 和 Stream Reactor Connectors。

我正在使用 Kafka 0.10.0.1

我正在使用 DataMountaineer Stream Reactor 0.2.4

我将 Stream Reactor 的 jar 文件放到了 Kafka libs 文件夹中,并以分布式模式运行 Kafka Connect

bin/connect-distributed.sh config/connect-distributed.properties

我添加了 Cassandra Source 连接器,如下所示:

curl -X POST -H "Content-Type: application/json" -d @config/connect-idoc-cassandra-source.json.txt localhost:8083/connectors

当我将数据添加到 Cassandra 表时,我看到它正在使用 Kafka 命令行使用者添加到主题中

bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic idocs-topic --from-beginning

以下是目前正在写入主题的示例:

{
"schema": {
    "type": "struct",
    "fields": [{
        "type": "string",
        "optional": true,
        "field": "idoc_id"
    }, {
        "type": "string",
        "optional": true,
        "field": "idoc_event_ts"
    }, {
        "type": "string",
        "optional": true,
        "field": "json_doc"
    }],
    "optional": false,
    "name": "idoc.idocs_events"
},
"payload": {
    "idoc_id": "dc4ab8a0-fdf8-11e6-8285-1bce55915fdd",
    "idoc_event_ts": "dc4ab8a1-fdf8-11e6-8285-1bce55915fdd",
    "json_doc": "{\"foo\":\"bar\"}"
}}

我想写给这个主题的是json_doc 列的值。

这是我的 Cassandra 源代码配置中的内容

{
"name": "cassandra-idocs",
"config": {
    "tasks.max": "1",
    "connector.class": "com.datamountaineer.streamreactor.connect.cassandra.source.CassandraSourceConnector",
    "connect.cassandra.key.space": "idoc",
    "connect.cassandra.source.kcql": "INSERT INTO idocs-topic SELECT json_doc FROM idocs_events PK idoc_event_ts",
    "connect.cassandra.import.mode": "incremental",
    "connect.cassandra.contact.points": "localhost",
    "connect.cassandra.port": 9042,
    "connect.cassandra.import.poll.interval": 10000
}}

如何更改 Kafka Connect Cassandra Source 的配置方式,以便仅将 json_doc 的值写入主题,使其看起来像这样:

{"foo":"bar"}

Kassandra Connect Query Language 似乎是可行的方法,但它并不限制写入 KCQL 中指定的列的内容。

更新

看到这个answer on StackOverflow 并将connect-distributed.properties 文件中的转换器从JsonConverter 更改为StringConverter

结果是现在写到了主题中:

Struct{idoc_id=74597cf0-fdf7-11e6-8285-1bce55915fdd,idoc_event_ts=74597cf1-fdf7-11e6-8285-1bce55915fdd,json_doc={"foo":"bar"}}

更新 2

connect-distributed.properties 文件中的转换器改回JsonConverter。然后还禁用了架构。

key.converter.schemas.enable=false
value.converter.schemas.enable=false 

结果是现在写到了主题中:

{
    "idoc_id": "dc4ab8a0-fdf8-11e6-8285-1bce55915fdd",
    "idoc_event_ts": "dc4ab8a1-fdf8-11e6-8285-1bce55915fdd",
    "json_doc": "{\"foo\":\"bar\"}"
}

注意 使用快照版本中的代码并将 KCQL 更改为

INSERT INTO idocs-topic 
SELECT json_doc, idoc_event_ts
FROM idocs_events 
IGNORE idoc_event_ts
PK idoc_event_ts

在主题上产生此结果

{"json_doc": "{\"foo\":\"bar\"}"}

谢谢

【问题讨论】:

  • 你得到了正确的答案。很高兴看到你想通了:)
  • 这仍然不是我们想要的。理想情况下,它只是 json_doc 列的值。还提交了两个 PR,以便 KCQL 在制作响应时使用 SELECT 和 IGNORE。我会将 UPDATE 2 中的结果更改为与 0.2.4 一起使用的方式

标签: cassandra apache-kafka apache-kafka-connect


【解决方案1】:

事实证明,在 DataMountaineer Stream Reactor 0.2.4 的 Cassandra 源中,我试图做的事情是不可能的。但是,快照版本(我假设将成为版本 0.2.5)将支持这一点。

下面是它的工作原理:

1) 将connect-distributed.properties文件中的转换器设置为StringConverter

2) 将 Cassandra 源连接器的 JSON 配置中的 KCQL 设置为

INSERT INTO idocs-topic 
SELECT json_doc, idoc_event_ts 
  FROM idocs_events 
IGNORE idoc_event_ts 
    PK idoc_event_ts 
WITHUNWRAP

这将导致json_doc 列的值在没有任何架构信息或列名本身的情况下发布到 Kafka 主题。

因此,如果列 json_doc 包含值 {"foo":"bar"},那么主题上会出现以下内容:

{"foo":"bar"}

以下是有关 KCQL 在快照版本中如何工作的一些背景信息。

SELECT 现在将只检索该表中在 KCQL 中指定的列。最初它总是检索所有的列。需要注意的是,在使用incremental 导入模式时,PK 列必须是SELECT 语句的一部分。如果 PK 列的值不应该包含在发布到 Kafka 主题的消息中,则将其添加到 IGNORE 语句中(如上例所示)。

WITHUNWRAP 是 KCQL 的新功能,它将告诉 Cassandra 源连接器使用 String Schema 类型(而不是 Struct)创建一个 SourceRecord。在这种模式下,只有SELECT 语句中的列的值将被存储为SourceRecord 的值。如果在应用IGNORE 语句后SELECT 语句中有多个列,则将这些值附加在一起并用逗号分隔。

【讨论】:

    猜你喜欢
    • 2017-05-27
    • 2019-04-23
    • 2020-10-25
    • 2019-08-16
    • 2019-04-05
    • 2021-01-30
    • 2017-01-08
    • 2018-10-16
    • 2016-12-02
    相关资源
    最近更新 更多