【发布时间】: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