集成

Apache Kafka

qianmoQqianmoQ· 更新于 2026-10-01· 阅读 16 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

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.properties

3) 启动 Kafka:

bin/kafka-server-start.sh config/server.properties

4) 创建 "test" 主题:

bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic test

5) 启动消费者:

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/connectors

Sink 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。

评论

登录后参与评论

正在加载评论…