工具

事件泵

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

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

虽然 Apache PLC4X API 让你能够简便地访问 PLC 资源,但反复读取同一组值这类工作,纯 API 是交给你自己来完成的:安排读取任务、租借与归还连接、决定 PLC 停止响应时该怎么办,以及确保缓慢的响应不会堆积在下一次轮询之后。

Event-Pump 会替你完成这些工作。

你只需描述读取什么以及何时读取,它就会把每个响应投递给一个监听器。

Event-Pump 取代了 PLC4X 0.13 及更早版本中的 Scraper。如果你正在迁移,请参阅本页末尾的从 Scraper 迁移。

Event-Pump 快速上手

Event-Pump 位于以下 Maven 模块中:

    <dependency>
      <groupId>org.apache.plc4x</groupId>
      <artifactId>plc4j-tools-event-pump</artifactId>
      <version>1.0.0</version>
    </dependency>

概念

事件泵(Event-Pump)由三部分组成:

TagBatch

一组标签从同一个连接读取,共用一个触发器和一个监听器。批(Batch)是调度的最小单位:其中的每个标签都在单次请求中被读取。

Trigger

决定何时获取一批数据。TimerTrigger 以固定间隔触发。

TagBatchListener

接收每个 PlcReadResponse 以及任何错误。

EventPump

一个容器,持有若干批次(batch),并同时启动和停止它们。

在 Java 中使用

批次通过构建器(builder)组装后交给泵(pump):

PlcConnectionFactory connectionFactory = PlcDriverManager.getDefault().getConnectionFactory();

TagBatch batch = TagBatch.builder()
    .withBatchId("boiler")
    .withConnectionFactory(connectionFactory)
    .withConnectionString("opcua:tcp://192.168.1.1:4840?request-timeout-ms=10000")
    .addTagAddress("temperature", "ns=2;i=1001")
    .addTagAddress("pressure", "ns=2;i=1002")
    .withTrigger(new TimerTrigger(5, TimeUnit.SECONDS))
    .withListener((b, response) -> {
        System.out.println("Temperature: " + response.getInt("temperature"));
        System.out.println("Pressure: " + response.getInt("pressure"));
    })
    .build();

EventPump pump = new EventPump();
pump.addBatch(batch);
pump.startAll();

// ... later
pump.close();

EventPump 实现了 AutoCloseable,关闭它会停止并关闭它所拥有的每个批次,因此 try-with-resources 代码块同样适用。

单个批次可以通过 id 使用 startBatch(id)、stopBatch(id) 和 removeBatch(id) 进行控制。

正在运行的批次的标签集可以随时通过 addTag、removeTag 和 clearTags 进行修改;该修改将在下一个抓取周期生效。

处理错误

将 lambda 作为监听器传入只能覆盖成功的情况。实现该接口即可同时看到失败情况:

batch.setListener(new TagBatch.TagBatchListener() {
    @Override
    public void onTagsFetched(TagBatch batch, PlcReadResponse response) {
        // process values
    }

    @Override
    public void onError(TagBatch batch, Throwable error) {
        // connection lost, read failed, ...
    }

    @Override
    public void onFetchSkipped(TagBatch batch, long lastFetchDurationMs, long consecutiveSkips) {
        // the PLC is slower than the configured interval
    }
});

配置文件

除了在代码中组装批次之外,整个泵机也可以用 YAML、JSON 或 XML 文件来描述,并通过 EventPumpFactory 加载:

connections:
  - id: plc1
    url: "opcua:tcp://192.168.1.1:4840?request-timeout-ms=10000"

batches:
  - id: boiler
    connectionId: plc1
    tags:
      temperature: "ns=2;i=1001"
      pressure: "ns=2;i=1002"
    trigger:
      type: timer
      intervalSeconds: 5
EventPump pump = EventPumpFactory.fromYaml(
    new File("event-pump.yml"), connectionManager, listener);
pump.startAll();

EventPumpFactory 还提供了 fromJson(…​) 和 fromXml(…​),且每个变体都接受一个可选的 ValueTransformerRegistry 作为最后一个参数。

传递给工厂的监听器将成为文件中所有批次的默认监听器。

触发间隔

定时器触发的间隔可以以秒为单位给出,也可以——对于 PLC 数据采集中常见的亚秒级频率——以毫秒为单位给出:

    trigger:
      type: timer
      intervalMillis: 250
      initialDelayMillis: 100

intervalSeconds/initialDelaySeconds 与 intervalMillis/initialDelayMillis 对于同一项设置是互斥的:如果在两种单位中都给出了同一项设置,程序会在启动时直接报错,而不是悄悄地选用其中一个,因为否则读取该文件的人就必须去猜测实际的速率。对一项设置使用秒、对另一项设置使用毫秒则是允许的。

转换

标签可以携带一个表达式,该表达式会在监听器看到其值之前对值进行处理。使用扩展标签格式来声明它:

batches:
  - id: boiler
    connectionId: plc1
    tags:
      temperature:
        address: "ns=2;i=1001"
        transform: "value * 1.8 + 32"
      pressure:
        address: "ns=2;i=1002"
    trigger:
      type: timer
      intervalSeconds: 5

在表达式中,value 指代标签自身的值,而同一批次中的其他任何标签都可以通过其自身名称引用,因此像 temperature + humidity 这样的跨标签表达式也能正常工作。所有名称都解析为当前响应中的值,且在任何转换应用之前即已完成解析。

内置的求值器(注册名称为 simple)支持:

  • 算术运算:+、-、*、/、%、一元负号和括号
  • 比较运算:>、<、>=、⇐、==、!=
  • 布尔逻辑:&&、||、!,以及字面量 true 和 false

如果表达式求值失败,系统会记录该错误,并将原始值原样传递,因此一个出错的表达式只会影响单个标签,而不会导致整个批次失败。

自定义转换器可以通过实现 ValueTransformer 并将其添加到 ValueTransformerRegistry 中来注册。

通过构建器也能使用同样的功能,即调用 addTransform(tagName, expression)。

超时、退避与过载

这些配置的正确性比 API 本身更为重要,因此有必要明确说明每个调节项各自的作用。

请求超时应配置在连接串上

Event-Pump 本身不设置请求超时。一次读取允许花费多长时间由驱动决定,可作为参数配置在连接 URL 上,例如:

opcua:tcp://192.168.1.1:4840?request-timeout-ms=10000

请查阅所用驱动的文档,了解其支持的参数。

抓取看门狗

作为最后一道防线,针对某个驱动始终无法完成读取的情况,每个批次都会为单次抓取周期设定上限。默认值为 5 分钟。这个值刻意远高于任何合理的请求超时:它并非用于限制一次读取可以耗时多久,而只是保证一个批次不会因为停滞的请求而永久卡住。

在构建器上使用 withFetchTimeout(long, TimeUnit) 设置,或在配置文件中的批次上使用 fetchTimeoutMs 设置;值为 0 或更小时将禁用该看门狗。

不要将看门狗当作驱动请求超时的替代品。将其设置为低于驱动超时的值,会使每一次缓慢的读取都被视为停滞的读取。

当 PLC 无法访问时

抓取失败后,批次会在再次尝试前进行指数退避:1 秒,然后是 2、4、8……上限为 60 秒。第一次成功的抓取会将其重置。这样可以避免批次反复冲击宕机的设备,也能让故障期间的日志保持可读。两端均可通过 withInitialBackoffMs(…​) 和 withMaxBackoffMs(…​) 进行配置。

当 PLC 慢于采集间隔时

如果上一次抓取仍在运行时触发器再次触发,新的抓取将被跳过而非排队,同时会记录一条警告,并在监听器上调用 onFetchSkipped。这是轮询间隔对该设备设置得过于激进的主要信号。

TimerTrigger 采用固定延迟而非固定速率进行调度,因此偶尔较慢的周期只会推迟下一次抓取,而不会产生一批追赶式读取的突发。

线程模型

每个 TimerTrigger 拥有一个定时器线程和一个分发线程;监听器代码在分发线程上运行,绝不会在定时器线程上运行。因此,阻塞的监听器只会延迟它自己所属的批次。

请注意,这意味着同一批次的监听器回调是串行执行的,而不同批次的回调可能并发运行——在多个批次之间共享的监听器必须是线程安全的。

Timer 也可以通过传入 TimerTrigger 构造函数在多个触发器之间共享,当批次数量很多时,这有助于减少线程数量。

从 Scraper 迁移

Scraper 已在 PLC4X 0.13 之后被移除。相关概念可以较为直接地对应过来:

ScraperEvent-Pump
ScrapeJobTagBatch
TriggeredScraperImplEventPump
ResultHandlerTagBatch.TagBatchListener
scrape rateTimerTrigger interval
futureTimeOut 构造函数参数驱动自身的 request-timeout-ms 连接字符串参数

最后一行值得关注。Scraper 无论连接字符串怎么写,都对每次读取施加自己统一的超时——默认 2000 毫秒,除非你另行指定。而 Event-Pump 并非如此:在连接 URL 上配置超时,驱动就会遵守该设置。

当前限制

  • SubscriptionTrigger 是一个占位符。它可以被构建,配置文件中也可以引用 type: subscription,但启动这样的批处理会抛出 UnsupportedOperationException。在订阅支持落地之前,请使用 TimerTrigger。
  • 批处理只支持读取,没有写入支持。
  • 为一次抓取而租用连接是一个在批处理调度线程上执行的阻塞调用。

评论

登录后参与评论

正在加载评论…