Apache Kafka
Apache Kafka 是一个开源的分布式事件流平台,被数千家公司用于构建高性能数据管道、流式分析、数据集成以及关键业务应用。
PLC4X Kafka 连接器
PLC4X 连接器能够在 Kafka 与使用工业协议的设备之间传递数据。它们可以从最新的 plc4x extras - kafka integration 以及 PLC4X 的 plc4x release 源码构建,也可以从 github 获取最新的快照构建。
简介
Connect 工作进程本质上是一个生产者或消费者进程,它提供一套标准 API,供 Kafka 用来管理它。它可以以两种模式运行:
- 独立(Standalone)
- 分布式(Distributed)
独立模式允许你从命令行在本地运行连接器,无需将 jar 文件安装到 Kafka 代理上。在分布式模式下,连接器运行在 Kafka 代理上,这需要你将 jar 文件安装到所有代理上。该模式允许工作进程分布在多个 Kafka 代理上,从而提供冗余和负载均衡。
快速开始
要启动 Kafka Connect 系统,需要执行以下步骤:
1) 从这里下载最新版本的 Apache Kafka 二进制包:https://kafka.apache.org/downloads。
2) 解压该归档文件。
3) 将 target/plc4j-apache-kafka-1.0.0-uber-jar.jar 复制到 Kafka 的 libs 目录,或复制到 config/connect-distributed.properties 文件中指定的插件目录。
4) 将 config 中的文件复制到 Kafka 的 config 目录。
5) 确保 OPCUA 服务器在发现过程中公布的主机名能够从 Kafka Connect 服务器解析到。最简单的方法是将该主机名添加到你的 hosts 文件中。
启动 Kafka 代理
1) 打开 4 个控制台窗口,并切换到该目录 2) 启动 Zookeeper:
bin/zookeeper-server-start.sh config/zookeeper.properties3) 启动 Kafka:
bin/kafka-server-start.sh config/server.properties4) 创建 "test" 主题:
bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic test5) 启动消费者:
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning源连接器
连接工作进程(connect worker)的初始配置由配置文件提供。不过,一旦工作进程启动,就可以通过 connect REST API 来修改配置,该 API 通常在 http://localhost:8083/connectors] 上可用。在分布式模式下运行时,所有配置都需要通过 REST API 来完成。
PLC4X Kafka 集成仓库的 config/plc4x-source.properties 目录中提供了一个示例配置文件。该文件包含注释以及可供工作进程使用的有效属性。
通过 REST 接口配置连接器时,需要提供与示例 config/plc4x-source.properties] 文件中指定的相同属性。这些属性需要以 JSON 格式提供,并附带若干请求头。下面的示例展示了所期望的格式:
curl -X POST -H "Content-Type: application/json" --data '{"name": "plc-source-test", "config": {"connector.class":"org.apache.plc4x.kafka.Plc4xSourceConnector",
// TODO: Continue here ...
"tasks.max":"1", "file":"test.sink.txt", "topics":"connect-test" }}' http://localhost:8083/connectors启动 Kafka Connect 源工作进程(独立模式)
适用于测试。
1) 启动 Kafka Connect:
bin/connect-standalone.sh config/connect-standalone.properties config/plc4x-source.properties现在通过 kafka-console-consumer 观察控制台窗口。
如果要调试连接器,请务必在启动 Kafka-Connect 之前设置一些环境变量:
export KAFKA_DEBUG=y; export DEBUG_SUSPEND_FLAG=y;在这种情况下,启动过程将暂停,直到有 IDE 通过远程调试会话连接为止。
启动 Kafka Connect 源工作进程(分布式模式)
适用于生产环境。
在这种情况下,节点的状态由 Zookeeper 管理,连接器的配置通过 Kafka 主题进行分发。
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-configs --replication-factor 3 --partitions 1 --config cleanup.policy=compact
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-offsets --replication-factor 3 --partitions 50 --config cleanup.policy=compact
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-status --replication-factor 3 --partitions 10 --config cleanup.policy=compact启动 Worker 也就如此简单:
bin /connect-distributed.sh config/connect-distributed.properties然后通过 REST 接口提供连接器的配置:
curl -X POST -H "Content-Type: application/json" --data '{"name": "plc-source-test", "config": {"connector.class":"org.apache.plc4x.kafka.Plc4xSourceConnector",
// TODO: Continue here ...
"tasks.max":"1", "file":"test.sink.txt", "topics":"connect-test" }}' http://localhost:8083/connectorsSink Connector
参见 config/sink.properties 获取示例配置。
启动 Kafka Connect Sink 工作进程(Standalone 模式)
非常适合用于测试。
1) 启动 Kafka Connect:
bin/connect-standalone.sh config/connect-standalone.properties config/plc4x-sink.properties现在使用 "kafka-console-producer" 打开控制台窗口。
使用下面示例数据包向 Kafka 主题生产数据,应当会使负载中包含的所有值按照 sink 属性中定义的映射关系发送到 PLC。
{"schema":
{"type":"struct","fields":
[{"type":"struct","fields":
[{"type":"boolean","optional":true,"field":"running"},
{"type":"boolean","optional":true,"field":"conveyorLeft"},
{"type":"boolean","optional":true,"field":"conveyorRight"},
{"type":"boolean","optional":true,"field":"load"},
{"type":"int32","optional":true,"field":"numLargeBoxes"},
{"type":"boolean","optional":true,"field":"unload"},
{"type":"boolean","optional":true,"field":"transferRight"},
{"type":"boolean","optional":true,"field":"transferLeft"},
{"type":"boolean","optional":true,"field":"conveyorEntry"},
{"type":"int32","optional":true,"field":"numSmallBoxes"}],
"optional":false,"name":"org.apache.plc4x.kafka.schema.Field","field":"fields"},
{"type":"int64","optional":false,"field":"timestamp"},
{"type":"int64","optional":true,"field":"expires"}],
"optional":false,"name":"org.apache.plc4x.kafka.schema.JobResult",
"doc":"PLC Job result. This contains all of the received PLCValues as well as a recieved timestamp"},
"payload":
{"fields":
{"running":false,"conveyorLeft":true,
"conveyorRight":true,"load":false,
"numLargeBoxes":1630806456,
"unload":true,
"transferRight":false,
"transferLeft":true,
"conveyorEntry":false,
"numSmallBoxes":-1135309911},
"timestamp":1606047842350,
"expires":null}}如果你想调试该连接器,请务必在启动 Kafka-Connect 之前设置一些环境变量:
export KAFKA_DEBUG=y; export DEBUG_SUSPEND_FLAG=y;在这种情况下,启动过程将会挂起,直到有 IDE 通过远程调试会话连接为止。
启动 Kafka Connect Sink 工作进程(分布式模式)
这种方式适用于生产环境。
在这种情况下,节点的状态由 Zookeeper 管理,连接器的配置通过 Kafka 主题进行分发。
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-configs --replication-factor 3 --partitions 1 --config cleanup.policy=compact
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-offsets --replication-factor 3 --partitions 50 --config cleanup.policy=compact
bin/kafka-topics --create --zookeeper localhost:2181 --topic connect-status --replication-factor 3 --partitions 10 --config cleanup.policy=compact启动工作者就像这样简单:
bin /connect-distributed.sh config/connect-distributed.properties随后,连接器的配置通过 REST 接口提供:
curl -X POST -H "Content-Type: application/json" --data '{"name": "plc-sink-test", "config": {"connector.class":"org.apache.plc4x.kafka.Plc4xSinkConnector",
// TODO: Continue here ...
"tasks.max":"1", "file":"test.sink.txt", "topics":"connect-test" }}' http://localhost:8083/connectors优雅退避
当读取或写入 PLC 地址发生错误时,系统实现了优雅退避机制,以避免向 PLC 发送大量请求。不过,由于每个 PLC 的连接器数量应当加以限制以减轻 PLC 的负载,因此优雅退避不应对性能产生重大影响。
对于 source 连接器,PLC4X 的抓取逻辑能够在发生故障时处理随机轮询速率,这些速率在连接器内部进行缓冲,连接器的轮询速率不会影响 PLC 的轮询速率。
对于 sink 连接器,如果写入失败,系统会以可配置的重试次数进行重试,每次重试之间存在超时间隔。系统会抛出一个可重试异常(Retriable Exception),为重试时机提供抖动(jitter)。
Schema 兼容性
PLC4X 定义了一个非常基础的 schema,大部分实现交由用户完成。它包含以下字段:
"fields":这是一个自定义结构,由连接器配置中定义的字段构成。它允许用户在此定义任意字段,且所有字段均基于 PLC4X 数据类型。"timestamp":这是 PLC4X 连接器处理该 PLC 请求时的时间戳。"expires":该字段由 sink 连接器使用,允许其丢弃过旧的记录。值为 0 或 null 表示无论记录多旧都不会被丢弃。
由于 schema 的大部分内容由用户定义,我们期望能够在基础 schema 之间提供向后兼容性。
sink 连接器和 source 连接器的 schema 是相同的。这使我们能够从一个 PLC 生产数据,并将数据发送到 sink。
评论
登录后参与评论
KnowForge