事件泵
虽然 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: 5EventPump pump = EventPumpFactory.fromYaml(
new File("event-pump.yml"), connectionManager, listener);
pump.startAll();EventPumpFactory 还提供了 fromJson(…) 和 fromXml(…),且每个变体都接受一个可选的 ValueTransformerRegistry 作为最后一个参数。
传递给工厂的监听器将成为文件中所有批次的默认监听器。
触发间隔
定时器触发的间隔可以以秒为单位给出,也可以——对于 PLC 数据采集中常见的亚秒级频率——以毫秒为单位给出:
trigger:
type: timer
intervalMillis: 250
initialDelayMillis: 100intervalSeconds/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 之后被移除。相关概念可以较为直接地对应过来:
| Scraper | Event-Pump |
|---|---|
ScrapeJob | TagBatch |
TriggeredScraperImpl | EventPump |
ResultHandler | TagBatch.TagBatchListener |
| scrape rate | TimerTrigger interval |
futureTimeOut 构造函数参数 | 驱动自身的 request-timeout-ms 连接字符串参数 |
最后一行值得关注。Scraper 无论连接字符串怎么写,都对每次读取施加自己统一的超时——默认 2000 毫秒,除非你另行指定。而 Event-Pump 并非如此:在连接 URL 上配置超时,驱动就会遵守该设置。
当前限制
SubscriptionTrigger是一个占位符。它可以被构建,配置文件中也可以引用type: subscription,但启动这样的批处理会抛出UnsupportedOperationException。在订阅支持落地之前,请使用TimerTrigger。- 批处理只支持读取,没有写入支持。
- 为一次抓取而租用连接是一个在批处理调度线程上执行的阻塞调用。
评论
登录后参与评论
KnowForge