你不会是程序猿吧?

当Elasticsearch遇见Kafka--Kafka Connect

mushao999 发表了文章 • 2 个评论 • 14349 次浏览 • 2018-11-17 11:15 • 来自相关话题

本文同步发布在腾讯云+社区Elasticsearch专栏中:https://cloud.tencent.com/developer/column/4008
在“[当Elasticsearch遇见Kafka--Logstash kafka input插件](https://cloud.tencent.com/developer/article/1362320)”一文中,我对Logstash的Kafka input插件进行了简单的介绍,并通过实际操作的方式,为大家呈现了使用该方式实现Kafka与Elastisearch整合的基本过程。可以看出使用Logstash input插件的方式,具有配置简单,数据处理方便等优点。然而使用Logstash Kafka插件并不是Kafka与Elsticsearch整合的唯一方案,另一种比较常见的方案是使用Kafka的开源组件Kafka Connect。

![Confluent实现Kafka与Elasticsearch的连接](https://main.qcloudimg.com/raw ... a3.png)

1 Kafka Connect简介


Kafka Connect是Kafka的开源组件Confluent提供的功能,用于实现Kafka与外部系统的连接。Kafka Connect同时支持分布式模式和单机模式,另外提供了一套完整的REST接口,用于查看和管理Kafka Connectors,还具有offset自动管理,可扩展等优点。

Kafka connect分为企业版和开源版,企业版在开源版的基础之上提供了监控,负载均衡,副本等功能,实际生产环境中建议使用企业版。(本测试使用开源版)

Kafka connect workers有两种工作模式,单机模式和分布式模式。在开发和适合使用单机模式的场景下,可以使用standalone模式, 在实际生产环境下由于单个worker的数据压力会比较大,distributed模式对负载均和和扩展性方面会有很大帮助。(本测试使用standalone模式)

关于Kafka Connect的详细情况可以参考[[Kafka Connect](https://docs.confluent.io/curr ... x.html)]

2 使用Kafka Connect连接Kafka和Elasticsearch


2.1 测试环境准备


本文与使用Logstash Kafka input插件环境一样[[传送门](https://cloud.tencent.com/developer/article/1362320)],组件列表如下

| 服务 | ip | port |
|:----|:----|:----|
| Elasticsearch service | 192.168.0.8 | 9200 |
| Ckafka | 192.168.13.10 | 9092 |
| CVM | 192.168.0.13 | - |


kafka topic也复用原来了的kafka_es_test

2.2 Kafka Connect 安装


[[Kafka Connec下载地址](https://www.confluent.io/download/)]

本文下载的为开源版本confluent-oss-5.0.1-2.11.tar.gz,下载后解压

2.3 Worker配置


1) 配置参考

如前文所说,worker分为Standalone和Distributed两种模式,针对两种模式的配置,参考如下

[[通用配置](https://www.confluent.io/download/)]

[[Standalone Woker配置](https://www.confluent.io/download/)]

[[Distributed Worker配置](https://www.confluent.io/download/)]

此处需要注意的是Kafka Connect默认使用AvroConverter,使用该AvroConverter时需要注意必须启动Schema Registry服务

2) 实际操作

本测试使用standalone模式,因此修改/root/confluent-5.0.1/etc/schema-registry/connect-avro-standalone.properties

<br /> bootstrap.servers=192.168.13.10:9092<br />

2.4 Elasticsearch Connector配置


1) 配置参考

[[Connectors通用配置](https://docs.confluent.io/curr ... ectors)]

[[Elasticsearch Configuration Options](https://docs.confluent.io/curr ... s.html)]

2) 实际操作

修改/root/confluent-5.0.1/etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties

<br /> name=elasticsearch-sink<br /> connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector<br /> tasks.max=1<br /> topics=kafka_es_test<br /> key.ignore=true<br /> connection.url=http://192.168.0.8:9200<br /> type.name=kafka-connect<br />

注意: 其中topics不仅对应Kafka的topic名称,同时也是Elasticsearch的索引名,当然也可以通过topic.index.map来设置从topic名到Elasticsearch索引名的映射

2.5 启动connector


1 注意事项

1) 由于配置文件中jar包位置均采用的相对路径,因此建议在confluent根目录下执行命令和启动程序,以避免不必要的问题

2) 如果前面没有修改converter,仍采用AvroConverter, 注意需要在启动connertor前启动Schema Registry服务

2 启动Schema Registry服务

正如前文所说,由于在配置worker时指定使用了AvroConverter,因此需要启动Schema Registry服务。而该服务需要指定一个zookeeper地址或Kafka地址,以存储schema数据。由于CKafka不支持用户通过接口形式创建topic,因此需要在本机起一个kafka以创建名为_schema的topic。

1) 启动Zookeeper

<br /> ./bin/zookeeper-server-start -daemon etc/kafka/zookeeper.properties<br />

2) 启动kafka

<br /> ./bin/kafka-server-start -daemon etc/kafka/server.properties<br />

3) 启动schema Registry

<br /> ./bin/schema-registry-start -daemon etc/schema-registry/schema-registry.properties<br />

4) 使用netstat -natpl 查看各服务端口是否正常启动

zookeeper 2181

kafka 9092

schema registry 8081

3 启动Connector

<br /> ./bin/connect-standalone -daemon etc/schema-registry/connect-avro-standalone.properties etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties<br />

ps:以上启动各服务均可在logs目录下找到对应日志

2.6 启动Kafka Producer


由于我们采用的是AvroConverter,因此不能采用Kafka工具包中的producer。Kafka Connector bin目录下提供了Avro Producer

1) 启动Producer

<br /> ./bin/kafka-avro-console-producer --broker-list 192.168.13.10:9092 --topic kafka_es_test --property value.schema='{"type":"record","name":"person","fields":[{"name":"nickname","type":"string"}]}'<br />

2) 输入如下数据

<br /> {"nickname":"michel"}<br /> {"nickname":"mushao"}<br />

2.7 Kibana验证结果


1) 查看索引

在kibana Dev Tools的Console中输入

<br /> GET _cat/indices<br />

结果

<br /> green open kafka_es_test 36QtDP6vQOG7ubOa161wGQ 5 1 1 0 7.9kb 3.9kb<br /> green open .kibana QUw45tN0SHqeHbF9-QVU6A 1 1 1 0 5.5kb 2.7kb<br />

可以看到名为kafka_es_test的索引被成功创建

2) 查看数据

<br /> {<br /> "took": 0,<br /> "timed_out": false,<br /> "_shards": {<br /> "total": 5,<br /> "successful": 5,<br /> "skipped": 0,<br /> "failed": 0<br /> },<br /> "hits": {<br /> "total": 2,<br /> "max_score": 1,<br /> "hits": [<br /> {<br /> "_index": "kafka_es_test",<br /> "_type": "kafka-connect",<br /> "_id": "kafka_es_test+0+0",<br /> "_score": 1,<br /> "_source": {<br /> "nickname": "michel"<br /> }<br /> },<br /> {<br /> "_index": "kafka_es_test",<br /> "_type": "kafka-connect",<br /> "_id": "kafka_es_test+0+1",<br /> "_score": 1,<br /> "_source": {<br /> "nickname": "mushao"<br /> }<br /> }<br /> ]<br /> }<br /> }<br />

可以看到数据已经被成功写入

3 Confluent CLI


3.1 简介


查阅资料时发现很多文章都是使用Confluent CLI启动Kafka Connect,然而官方文档已经明确说明了该CLI只是适用于开发阶段,不能用于生产环境。

它可以一键启动包括zookeeper,kafka,schema registry, kafka rest, connect等在内的多个服务。但是这些服务对于Kafka Connect都不是必须的,如果不使用AvroConverter,则只需要启动Connect即可。即使使用了AvroConverter, 也只需要启动schema registry,将schema保存在远端的kafka中。Kafka Connect REST API也只是为用户提供一个管理connector的接口,也不是必选的。

另外使用CLI启动默认配置为启动Distributed的Connector,需要通过环境变量来修改配置

3.2 使用Confluent CLI


confluent CLI提供了丰富的命令,包括服务启动,服务停止,状态查询,日志查看等,详情参考如下简介视频 [[Introducing the Confluent CLI | Screencast](https://www.youtube.com/watch?v=ZKqBptBHZTg)]

1) 启动

<br /> ./bin/confluent start<br />

2) 检查confluent运行状态

<br /> ./bin/confluent status<br />

当得到如下结果则说明confluent启动成功

<br /> ksql-server is [UP]<br /> connect is [UP]<br /> kafka-rest is [UP]<br /> schema-registry is [UP]<br /> kafka is [UP]<br /> zookeeper is [UP]<br />

3) 问题定位

如果第二步出现问题,可以使用log命令查看,如connect未启动成功则

<br /> ./bin/confluent log connect<br />

4) 加载Elasticsearch Connector

a) 查看connector

<br /> ./bin/confluent list connectors<br />

结果

<br /> Bundled Predefined Connectors (edit configuration under etc/):<br /> elasticsearch-sink<br /> file-source<br /> file-sink<br /> jdbc-source<br /> jdbc-sink<br /> hdfs-sink<br /> s3-sink<br />

b) 加载Elasticsearch connector

<br /> ./bin/confluent load elasticsearch-sink<br />

结果

<br /> {<br /> "name": "elasticsearch-sink",<br /> "config": {<br /> "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",<br /> "tasks.max": "1",<br /> "topics": "kafka_es_test",<br /> "key.ignore": "true",<br /> "connection.url": "<a href="http://192.168.0.8:9200"" rel="nofollow" target="_blank">http://192.168.0.8:9200"</a>,<br /> "type.name": "kafka-connect",<br /> "name": "elasticsearch-sink"<br /> },<br /> "tasks": [],<br /> "type": null<br /> }<br />

5) 使用producer生产数据,并使用kibana验证是否写入成功

4 Kafka Connect Rest API


Kafka Connect提供了一套完成的管理Connector的接口,详情参考[[Kafka Connect REST Interface](https://docs.confluent.io/curr ... i.html)]。该接口可以实现对Connector的创建,销毁,修改,查询等操作

1) GET connectors 获取运行中的connector列表

2) POST connectors 使用指定的名称和配置创建connector

3) GET connectors/(string:name) 获取connector的详细信息

4) GET connectors/(string:name)/config 获取connector的配置

5) PUT connectors/(string:name)/config 设置connector的配置

6) GET connectors/(string:name)/status 获取connector状态

7) POST connectors/(stirng:name)/restart 重启connector

8) PUT connectors/(string:name)/pause 暂停connector

9) PUT connectors/(string:name)/resume 恢复connector

10) DELETE connectors/(string:name)/ 删除connector

11) GET connectors/(string:name)/tasks 获取connectors任务列表

12) GET /connectors/(string: name)/tasks/(int: taskid)/status 获取任务状态

13) POST /connectors/(string: name)/tasks/(int: taskid)/restart 重启任务

14) GET /connector-plugins/ 获取已安装插件列表

15) PUT /connector-plugins/(string: name)/config/validate 验证配置

5 总结


Kafka Connect是Kafka一个功能强大的组件,为kafka提供了与外部系统连接的一套完整方案,包括数据传输,连接管理,监控,多副本等。相对于Logstash Kafka插件,功能更为全面,但配置也相对为复杂些。有文章提到其性能也优于Logstash Kafka Input插件,如果对写入性能比较敏感的场景,可以在实际压测的基础上进行选择。另外由于直接将数据从Kafka写入Elasticsearch, 如果需要对文档进行处理时,选择Logstash可能更为方便。

kibana 访问,或者es 本身 访问限制措施?

yayg2008 回复了问题 • 8 人关注 • 5 个回复 • 2937 次浏览 • 2018-11-19 16:45 • 来自相关话题

jest 和restClient 哪个好用点,有没有demo推荐?

rochy 回复了问题 • 4 人关注 • 3 个回复 • 3764 次浏览 • 2018-11-17 13:55 • 来自相关话题

elastic内存问题 yong survivor总量都是0

zhangg7723 回复了问题 • 2 人关注 • 2 个回复 • 2811 次浏览 • 2018-11-26 11:23 • 来自相关话题

es建立倒排索引时如何区分字段的?倒排索引是token对应文档的集合,里面是否区分token匹配到一篇文档的哪些字段?

weizijun 回复了问题 • 2 人关注 • 1 个回复 • 5845 次浏览 • 2018-11-16 12:46 • 来自相关话题

mysql数据导入es,多值字段如何实现

回复

haliaddel 回复了问题 • 1 人关注 • 1 个回复 • 1519 次浏览 • 2018-11-19 10:02 • 来自相关话题

es数组聚合

doctor 回复了问题 • 2 人关注 • 1 个回复 • 4583 次浏览 • 2018-11-20 17:37 • 来自相关话题

es multi_match 索引优化

rochy 回复了问题 • 2 人关注 • 1 个回复 • 2388 次浏览 • 2018-11-15 18:17 • 来自相关话题

关于ik_max_word分词的疑问

greenjim20 回复了问题 • 2 人关注 • 3 个回复 • 5970 次浏览 • 2018-11-16 18:55 • 来自相关话题

elastic from size 分页遇到重复数据问题

bznie 回复了问题 • 4 人关注 • 6 个回复 • 4534 次浏览 • 2022-07-06 10:47 • 来自相关话题

可以通过scroll_id对应到要查询的索引吗?

mushao999 回复了问题 • 5 人关注 • 4 个回复 • 2913 次浏览 • 2018-11-19 21:53 • 来自相关话题

请教一问题,elasticsearch如何实现句内检索或者段内检索

rochy 回复了问题 • 4 人关注 • 2 个回复 • 2402 次浏览 • 2018-11-15 18:13 • 来自相关话题

ElasticSearch中的id可以为中文嘛?

rochy 回复了问题 • 2 人关注 • 1 个回复 • 1517 次浏览 • 2018-11-15 11:54 • 来自相关话题

【 报名已结束】2018 Elastic & 东方航空大数据技术沙龙

kennywu76 发表了文章 • 0 个评论 • 2866 次浏览 • 2018-11-14 15:21 • 来自相关话题

本次活动报名已截止,因为名额限制无法报名成功的小伙伴也不用着急,届时会议将采用zoom进行直播,在 PC、Mac、iPhone/iPad、安卓手机/平板上,点击https://www.zoomus.cn/j/1524425455 即可轻松加入观看。
本次活动报名已截止,因为名额限制无法报名成功的小伙伴也不用着急,届时会议将采用zoom进行直播,在 PC、Mac、iPhone/iPad、安卓手机/平板上,点击https://www.zoomus.cn/j/1524425455 即可轻松加入观看。

search_template支持高亮吗?高亮无效

rochy 回复了问题 • 3 人关注 • 1 个回复 • 2646 次浏览 • 2018-11-14 15:36 • 来自相关话题