This module provides polling-based subscription emulation for drivers that don't natively support subscriptions.
The PollingSubscriptionConnectionBase class provides polling-based subscription emulation for drivers that don't natively support subscriptions (such as Modbus, BACnet, etc.).
UnsupportedOperationException)To add polling-based subscription support to your driver:
pom.xml:<dependency> <groupId>org.apache.plc4x</groupId> <artifactId>plc4j-utils-subscription-emulation</artifactId> <version>1.0.0-SNAPSHOT</version> </dependency>
PollingSubscriptionConnectionBase instead of ConnectionBase and implement PlcReader:import org.apache.plc4x.java.utils.subscriptionemulation.PollingSubscriptionConnectionBase; public class ModbusTcpConnection extends PollingSubscriptionConnectionBase<ModbusTcpConfiguration> implements PlcReader, PlcWriter { public ModbusTcpConnection(ModbusTcpConfiguration configuration, TransportInstance<?> transportInstance) { super(configuration, transportInstance); } // Implement PlcReader interface @Override public CompletableFuture<PlcReadResponse> read(PlcReadRequest readRequest) { // Your read implementation } // Rest of your driver implementation... }
That's it! Your driver now supports subscriptions via polling.
When a subscription is created:
read() methodYou can override these methods to customize behavior:
@Override protected long getDefaultPollingInterval() { return 500; // 500ms instead of default 1000ms }
@Override protected boolean valuesEqual(PlcValue v1, PlcValue v2) { // Custom comparison logic // For example, treat small floating point differences as equal if (v1 == v2) return true; if (v1 == null || v2 == null) return false; if (v1.isDouble() && v2.isDouble()) { double diff = Math.abs(v1.getDouble() - v2.getDouble()); return diff < 0.001; // Treat differences < 0.001 as equal } return Objects.equals(v1.getObject(), v2.getObject()); }
PlcSubscriptionRequest request = DefaultPlcSubscriptionRequest.builder() .addCyclicTagAddress("temperature", "MAIN.temperature", Duration.ofMillis(100)) .addCyclicTagAddress("pressure", "MAIN.pressure", Duration.ofMillis(100)) .setConsumer(event -> { System.out.println("Temperature: " + event.getInteger("temperature")); System.out.println("Pressure: " + event.getInteger("pressure")); }) .build(); CompletableFuture<PlcSubscriptionResponse> future = connection.subscribe(request); PlcSubscriptionResponse response = future.get(); // Later, to unsubscribe: PlcUnsubscriptionRequest unsubRequest = DefaultPlcUnsubscriptionRequest.builder() .addHandles(response.getSubscriptionHandle("temperature")) .addHandles(response.getSubscriptionHandle("pressure")) .build(); connection.unsubscribe(unsubRequest).get();
PlcSubscriptionRequest request = DefaultPlcSubscriptionRequest.builder() .addChangeOfStateTagAddress("alarm", "MAIN.alarm", Duration.ofMillis(50)) .setConsumer(event -> { System.out.println("Alarm state changed: " + event.getBoolean("alarm")); }) .build(); connection.subscribe(request);
PlcSubscriptionRequest request = DefaultPlcSubscriptionRequest.builder() .addCyclicTagAddress("tag1", "MAIN.tag1", Duration.ofMillis(100)) .setTagConsumer("tag1", event -> { System.out.println("Tag1 event: " + event.getInteger("tag1")); }) .addCyclicTagAddress("tag2", "MAIN.tag2", Duration.ofMillis(100)) .setTagConsumer("tag2", event -> { System.out.println("Tag2 event: " + event.getInteger("tag2")); }) .build(); connection.subscribe(request);
// Subscribe PlcSubscriptionResponse response = connection.subscribe(request).get(); PlcSubscriptionHandle handle1 = response.getSubscriptionHandle("tag1"); PlcSubscriptionHandle handle2 = response.getSubscriptionHandle("tag2"); // Register a consumer for specific handles PlcConsumerRegistration registration = connection.registerConsumer( event -> { System.out.println("Registered consumer received event"); }, Arrays.asList(handle1, handle2) ); // Later, unregister connection.unregisterConsumer(registration);
Your driver must:
PollingSubscriptionConnectionBase instead of ConnectionBasePlcReader interfaceplc4j-utils-subscription-emulation dependencyObjects.equals() by default, which may not be suitable for all data types (e.g., floating point numbers with small differences)For implementation details, see:
PollingSubscriptionConnectionBase.java:73 - Main implementationEventPump - The underlying polling engine (see java/tools/event-pump/README.md)