Chore/update asf yml (#757)
* refactor: Updated the kafka and calcite modules to use the new event-pump instead of the now obsolete scraper.
* chore: Updated the .asf.yml to match that of the main repo.
diff --git a/.asf.yaml b/.asf.yaml
index 4a1d4e7..322a356 100644
--- a/.asf.yaml
+++ b/.asf.yaml
@@ -76,14 +76,15 @@
excludes: []
bypass_teams:
- root
+ - plc4x-committers
restrict_deletion: true
restrict_force_push: true
required_pull_request_reviews:
dismiss_stale_reviews: true
- require_last_push_approval: false
+ require_last_push_approval: true
required_approving_review_count: 1
required_linear_history: false
- required_signatures: true
+ required_signatures: false
required_conversation_resolution: true
notifications:
diff --git a/plc4j/integrations/apache-calcite/pom.xml b/plc4j/integrations/apache-calcite/pom.xml
index 9b2b79a..c794fae 100644
--- a/plc4j/integrations/apache-calcite/pom.xml
+++ b/plc4j/integrations/apache-calcite/pom.xml
@@ -57,7 +57,7 @@
</dependency>
<dependency>
<groupId>org.apache.plc4x</groupId>
- <artifactId>plc4j-scraper</artifactId>
+ <artifactId>plc4j-tools-event-pump</artifactId>
<version>${plc4x.version}</version>
</dependency>
<dependency>
diff --git a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xBaseTable.java b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xBaseTable.java
index 625e4d3..445c185 100644
--- a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xBaseTable.java
+++ b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xBaseTable.java
@@ -29,7 +29,7 @@
import org.apache.calcite.schema.impl.AbstractTable;
import org.apache.calcite.util.ImmutableBitSet;
import org.apache.plc4x.java.api.exceptions.PlcRuntimeException;
-import org.apache.plc4x.java.scraper.config.JobConfiguration;
+import org.apache.plc4x.java.tools.eventpump.config.BatchConfiguration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -53,12 +53,12 @@
private static final Logger logger = LoggerFactory.getLogger(Plc4xBaseTable.class);
private final BlockingQueue<Plc4xSchema.Record> queue;
- private final JobConfiguration conf;
+ private final BatchConfiguration conf;
private final long tableCutoff;
private Plc4xSchema.Record current;
private final List<String> names;
- public Plc4xBaseTable(BlockingQueue<Plc4xSchema.Record> queue, JobConfiguration conf, long tableCutoff) {
+ public Plc4xBaseTable(BlockingQueue<Plc4xSchema.Record> queue, BatchConfiguration conf, long tableCutoff) {
this.tableCutoff = tableCutoff;
logger.info("Instantiating new PLC4X Table with configuration: {}", conf);
this.queue = queue;
@@ -107,7 +107,7 @@
} catch (ExecutionException | TimeoutException e) {
throw new PlcRuntimeException("Unable to fetch first record and infer arguments!", e);
}
- logger.info("Inferring types for Table '{}' based on values: {}", conf.getName(), first.values);
+ logger.info("Inferring types for Table '{}' based on values: {}", conf.getId(), first.values);
// Extract types
List<RelDataType> types = names.stream()
.map(n -> {
diff --git a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchema.java b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchema.java
index 9f1b3ee..6950a50 100644
--- a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchema.java
+++ b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchema.java
@@ -22,12 +22,12 @@
import org.apache.calcite.schema.impl.AbstractSchema;
import org.apache.plc4x.java.DefaultPlcDriverManager;
import org.apache.plc4x.java.api.exceptions.PlcRuntimeException;
-import org.apache.plc4x.java.scraper.ResultHandler;
-import org.apache.plc4x.java.scraper.Scraper;
-import org.apache.plc4x.java.scraper.ScraperImpl;
-import org.apache.plc4x.java.scraper.config.JobConfiguration;
-import org.apache.plc4x.java.scraper.config.ScraperConfiguration;
-import org.apache.plc4x.java.scraper.exception.ScraperException;
+import org.apache.plc4x.java.api.messages.PlcReadResponse;
+import org.apache.plc4x.java.tools.eventpump.EventPump;
+import org.apache.plc4x.java.tools.eventpump.TagBatch;
+import org.apache.plc4x.java.tools.eventpump.config.BatchConfiguration;
+import org.apache.plc4x.java.tools.eventpump.config.EventPumpConfiguration;
+import org.apache.plc4x.java.tools.eventpump.config.EventPumpFactory;
import org.apache.plc4x.java.utils.cache.CachedPlcConnectionManager;
import java.time.Instant;
@@ -38,36 +38,43 @@
public class Plc4xSchema extends AbstractSchema {
- protected final ScraperConfiguration configuration;
- protected final Scraper scraper;
+ protected final EventPumpConfiguration configuration;
+ protected final EventPump eventPump;
protected final QueueHandler handler;
protected final Map<String, BlockingQueue<Record>> queues;
protected final Map<String, Table> tableMap;
+ /** batch id -> connection id, so a record can be attributed to the PLC it came from. */
+ protected final Map<String, String> connectionIds;
- public Plc4xSchema(ScraperConfiguration configuration, long tableCutoff) throws ScraperException {
+ public Plc4xSchema(EventPumpConfiguration configuration, long tableCutoff) throws Exception {
this.configuration = configuration;
this.handler = new QueueHandler();
- this.scraper = new ScraperImpl(handler,
+ this.connectionIds = configuration.getBatches().stream()
+ .collect(Collectors.toMap(
+ BatchConfiguration::getId,
+ BatchConfiguration::getConnectionId
+ ));
+ this.queues = configuration.getBatches().stream()
+ .collect(Collectors.toMap(
+ BatchConfiguration::getId,
+ conf -> new ArrayBlockingQueue<>(1000)
+ ));
+ // Create the tables - one per batch
+ this.tableMap = configuration.getBatches().stream()
+ .collect(Collectors.toMap(
+ BatchConfiguration::getId,
+ conf -> defineTable(queues.get(conf.getId()), conf, tableCutoff)
+ ));
+ // Every batch reports to the same handler, which routes by batch id
+ this.eventPump = EventPumpFactory.create(configuration,
CachedPlcConnectionManager.getBuilder()
.withConnectionManager(new DefaultPlcDriverManager())
.build(),
- configuration.getJobs());
- this.queues = configuration.getJobConfigurations().stream()
- .collect(Collectors.toMap(
- JobConfiguration::getName,
- conf -> new ArrayBlockingQueue<>(1000)
- ));
- // Create the tables
- this.tableMap = configuration.getJobConfigurations().stream()
- .collect(Collectors.toMap(
- JobConfiguration::getName,
- conf -> defineTable(queues.get(conf.getName()), conf, tableCutoff)
- ));
- // Start the scraper
- this.scraper.start();
+ handler);
+ this.eventPump.startAll();
}
- Table defineTable(BlockingQueue<Record> queue, JobConfiguration configuration, Long limit) {
+ Table defineTable(BlockingQueue<Record> queue, BatchConfiguration configuration, Long limit) {
if (limit <= 0) {
return new Plc4xStreamTable(queue, configuration);
} else {
@@ -94,13 +101,16 @@
}
}
- class QueueHandler implements ResultHandler {
+ class QueueHandler implements TagBatch.TagBatchListener {
@Override
- public void handle(String job, String alias, Map<String, Object> results) {
+ public void onTagsFetched(TagBatch batch, PlcReadResponse response) {
+ String batchId = batch.getBatchId();
+ Map<String, Object> results = response.getTagNames().stream()
+ .collect(Collectors.toMap(name -> name, response::getObject));
try {
- Record record = new Record(Instant.now(), alias, results);
- queues.get(job).put(record);
+ Record record = new Record(Instant.now(), connectionIds.get(batchId), results);
+ queues.get(batchId).put(record);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new PlcRuntimeException("Handling got interrupted", e);
diff --git a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchemaFactory.java b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchemaFactory.java
index b287610..92e4d2e 100644
--- a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchemaFactory.java
+++ b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xSchemaFactory.java
@@ -22,12 +22,11 @@
import org.apache.calcite.schema.SchemaFactory;
import org.apache.calcite.schema.SchemaPlus;
import org.apache.commons.lang3.Validate;
-import org.apache.plc4x.java.scraper.config.ScraperConfiguration;
-import org.apache.plc4x.java.scraper.config.triggeredscraper.ScraperConfigurationTriggeredImpl;
-import org.apache.plc4x.java.scraper.exception.ScraperException;
+import org.apache.plc4x.java.tools.eventpump.config.EventPumpConfiguration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.io.File;
import java.io.IOException;
import java.util.Map;
@@ -40,9 +39,17 @@
Object config = operand.get("config");
Validate.notNull(config, "No configuration file given. Please specify operand 'config'...'");
// Load configuration from file
- ScraperConfiguration configuration;
+ EventPumpConfiguration configuration;
+ String configPath = config.toString();
try {
- configuration = ScraperConfiguration.fromFile(config.toString(), ScraperConfigurationTriggeredImpl.class);
+ File configFile = new File(configPath);
+ if (configPath.endsWith(".json")) {
+ configuration = EventPumpConfiguration.fromJson(configFile);
+ } else if (configPath.endsWith(".xml")) {
+ configuration = EventPumpConfiguration.fromXml(configFile);
+ } else {
+ configuration = EventPumpConfiguration.fromYaml(configFile);
+ }
} catch (IOException e) {
throw new IllegalArgumentException("Unable to load configuration file!", e);
}
@@ -59,8 +66,8 @@
// Pass the configuration to the Schema
try {
return new Plc4xSchema(configuration, parsedLimit);
- } catch (ScraperException e) {
- LOGGER.warn("Could not evaluate Plc4xSchema",e);
+ } catch (Exception e) {
+ LOGGER.warn("Could not evaluate Plc4xSchema", e);
//ToDo Exception, but interface does not accept ... null is fishy
return null;
}
diff --git a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xStreamTable.java b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xStreamTable.java
index 7c9668d..840c8ca 100644
--- a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xStreamTable.java
+++ b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xStreamTable.java
@@ -23,13 +23,13 @@
import org.apache.calcite.schema.ScannableTable;
import org.apache.calcite.schema.StreamableTable;
import org.apache.calcite.schema.Table;
-import org.apache.plc4x.java.scraper.config.JobConfiguration;
+import org.apache.plc4x.java.tools.eventpump.config.BatchConfiguration;
import java.util.concurrent.BlockingQueue;
public class Plc4xStreamTable extends Plc4xBaseTable implements ScannableTable, StreamableTable {
- public Plc4xStreamTable(BlockingQueue<Plc4xSchema.Record> queue, JobConfiguration conf) {
+ public Plc4xStreamTable(BlockingQueue<Plc4xSchema.Record> queue, BatchConfiguration conf) {
super(queue, conf, -1L);
}
diff --git a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xTable.java b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xTable.java
index 9c479a0..2b64b68 100644
--- a/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xTable.java
+++ b/plc4j/integrations/apache-calcite/src/main/java/org/apache/plc4x/Plc4xTable.java
@@ -21,13 +21,13 @@
import org.apache.calcite.DataContext;
import org.apache.calcite.linq4j.Enumerable;
import org.apache.calcite.schema.ScannableTable;
-import org.apache.plc4x.java.scraper.config.JobConfiguration;
+import org.apache.plc4x.java.tools.eventpump.config.BatchConfiguration;
import java.util.concurrent.BlockingQueue;
public class Plc4xTable extends Plc4xBaseTable implements ScannableTable {
- public Plc4xTable(BlockingQueue<Plc4xSchema.Record> queue, JobConfiguration conf, long tableCutoff) {
+ public Plc4xTable(BlockingQueue<Plc4xSchema.Record> queue, BatchConfiguration conf, long tableCutoff) {
super(queue, conf, tableCutoff);
}
diff --git a/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/DriverManagerTest.java b/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/DriverManagerTest.java
index c9f9323..8116acd 100644
--- a/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/DriverManagerTest.java
+++ b/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/DriverManagerTest.java
@@ -20,12 +20,11 @@
import org.apache.calcite.jdbc.CalciteConnection;
import org.apache.calcite.jdbc.Driver;
-import org.apache.plc4x.java.scraper.config.ScraperConfiguration;
-import org.apache.plc4x.java.scraper.config.ScraperConfigurationClassicImpl;
-import org.apache.plc4x.java.scraper.exception.ScraperException;
+import org.apache.plc4x.java.tools.eventpump.config.EventPumpConfiguration;
import org.assertj.core.api.WithAssertions;
import org.junit.jupiter.api.Test;
+import java.io.File;
import java.io.IOException;
import java.sql.Connection;
import java.sql.ResultSet;
@@ -36,12 +35,12 @@
public class DriverManagerTest implements WithAssertions {
@Test
- void query() throws SQLException, IOException, ScraperException {
+ void query() throws Exception {
Driver driver = new Driver();
Connection connection = driver.connect("jdbc:calcite:asdf;lex=MYSQL_ANSI", new Properties());
CalciteConnection calciteConnection = connection.unwrap(CalciteConnection.class);
- calciteConnection.getRootSchema().add("plc4x", new Plc4xSchema(ScraperConfiguration.fromFile("src/test/resources/example.yml", ScraperConfigurationClassicImpl.class), 100));
+ calciteConnection.getRootSchema().add("plc4x", new Plc4xSchema(EventPumpConfiguration.fromYaml(new File("src/test/resources/example.yml")), 100));
ResultSet rs = connection.prepareStatement("SELECT * FROM \"plc4x\".\"job1\"").executeQuery();
validateResult(rs);
diff --git a/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/Plc4XBaseTableTest.java b/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/Plc4XBaseTableTest.java
index 19c719a..932e7ae 100644
--- a/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/Plc4XBaseTableTest.java
+++ b/plc4j/integrations/apache-calcite/src/test/java/org/apache/plc4x/Plc4XBaseTableTest.java
@@ -19,7 +19,7 @@
package org.apache.plc4x;
import org.apache.calcite.linq4j.Enumerator;
-import org.apache.plc4x.java.scraper.config.JobConfigurationImpl;
+import org.apache.plc4x.java.tools.eventpump.config.BatchConfiguration;
import org.assertj.core.api.WithAssertions;
import org.junit.jupiter.api.Test;
@@ -33,12 +33,10 @@
@Test
void testOnBlockingQueue() {
ArrayBlockingQueue<Plc4xSchema.Record> queue = new ArrayBlockingQueue<>(100);
- Plc4xStreamTable table = new Plc4xStreamTable(queue, new JobConfigurationImpl(
- "job1",
- null,
- 100,
- Collections.emptyList(),
- Collections.singletonMap("key", "address")));
+ BatchConfiguration conf = new BatchConfiguration();
+ conf.setId("job1");
+ conf.setSimpleTags(Collections.singletonMap("key", "address"));
+ Plc4xStreamTable table = new Plc4xStreamTable(queue, conf);
Map<String, Object> objects = Collections.singletonMap("key", "value");
queue.add(new Plc4xSchema.Record(Instant.now(), "", objects));
diff --git a/plc4j/integrations/apache-calcite/src/test/resources/example.yml b/plc4j/integrations/apache-calcite/src/test/resources/example.yml
index d90dd3c..b0809a1 100644
--- a/plc4j/integrations/apache-calcite/src/test/resources/example.yml
+++ b/plc4j/integrations/apache-calcite/src/test/resources/example.yml
@@ -17,16 +17,28 @@
# under the License.
# ----------------------------------------------------------------------------
---
-sources:
- test: simulated:test
- test2: simulated:test2
+connections:
+ - id: test
+ url: simulated:test
+ - id: test2
+ url: simulated:test2
-jobs:
- - name: job1
- scrapeRate: 10
- sources:
- - test
- - test2
+# Every batch becomes one table, named after the batch id.
+batches:
+ - id: job1
+ connectionId: test
tags:
test: 'RANDOM/test:DINT'
test2: 'RANDOM/test:STRING'
+ trigger:
+ type: timer
+ intervalMillis: 10
+
+ - id: job2
+ connectionId: test2
+ tags:
+ test: 'RANDOM/test:DINT'
+ test2: 'RANDOM/test:STRING'
+ trigger:
+ type: timer
+ intervalMillis: 10
diff --git a/plc4j/integrations/apache-kafka/README.md b/plc4j/integrations/apache-kafka/README.md
index b8c5be7..13dbdbf 100644
--- a/plc4j/integrations/apache-kafka/README.md
+++ b/plc4j/integrations/apache-kafka/README.md
@@ -203,7 +203,8 @@
bombarded with requests. However as the number of connectors for each PLC should be limited to reduce the load on the PLC,
the graceful backoff shouldn't have a major impact.
-For the source connector the PLC4X scraper logic is able to handle randomized polling rates on failures, this is buffered within the
+For the source connector the PLC4X event-pump backs off exponentially on failures (1s, doubling up to 60s) and skips a poll
+if the previous read is still running, this is buffered within the
connector, the poll rate of the connector has no affect on the PLC poll rate.
For the sink connector, if a write fails it is retried a configurable number of times with a timeout between each time.
diff --git a/plc4j/integrations/apache-kafka/config/plc4x-source.properties b/plc4j/integrations/apache-kafka/config/plc4x-source.properties
index 6c024cb..13205b9 100644
--- a/plc4j/integrations/apache-kafka/config/plc4x-source.properties
+++ b/plc4j/integrations/apache-kafka/config/plc4x-source.properties
@@ -32,7 +32,7 @@
#This value controls how often it returns when no messages are received.
sources.machineA.pollReturnInterval=5000
-#There is an internal buffer between the PLC4X scraper and Kafka Connect.
+#There is an internal buffer between the PLC4X event-pump and Kafka Connect.
#This is the size of that buffer.
sources.machineA.bufferSize=1000
@@ -45,7 +45,7 @@
#A list of jobs specified in the following section.
jobs=simulated-dashboard,simulated-heartbeat
-#The poll rate for this job. the PLC4X scraper will request data every interval (ms).
+#The poll rate for this job. The PLC4X event-pump will request data every interval (ms).
jobs.simulated-dashboard.interval=1000
#A list of tags. Each tag is a map between an alias and a PLC4X address.
diff --git a/plc4j/integrations/apache-kafka/pom.xml b/plc4j/integrations/apache-kafka/pom.xml
index 18bc4bd..8d9ecc0 100644
--- a/plc4j/integrations/apache-kafka/pom.xml
+++ b/plc4j/integrations/apache-kafka/pom.xml
@@ -157,7 +157,7 @@
</dependency>
<dependency>
<groupId>org.apache.plc4x</groupId>
- <artifactId>plc4j-scraper</artifactId>
+ <artifactId>plc4j-tools-event-pump</artifactId>
<version>${plc4x.version}</version>
</dependency>
<dependency>
diff --git a/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceConnector.java b/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceConnector.java
index 5527075..86731a2 100644
--- a/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceConnector.java
+++ b/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceConnector.java
@@ -57,7 +57,7 @@
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// Initially we planned to have the simple assumption that one task maps to one PLC connection.
- // But we could easily say that one scraper instance maps to a task and one scraper task can
+ // But we could easily say that one event-pump instance maps to a task and one batch can
// process multiple PLC connections. But I guess this would be an optimization as we have to
// balance the load manually.
if(sourceConfig.getJobs().size() > maxTasks) {
@@ -66,8 +66,8 @@
return Collections.emptyList();
}
- // For each configured source we'll start a dedicated scraper instance collecting
- // all the scraper jobs enabled for this source.
+ // For each configured source we'll start a dedicated event-pump instance collecting
+ // all the jobs enabled for this source.
List<Map<String, String>> configs = new LinkedList<>();
for (Source source : sourceConfig.getSources()) {
// Build a list of job configurations only containing the ones referenced from
diff --git a/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceTask.java b/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceTask.java
index 9a4a353..3bc3f7a 100644
--- a/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceTask.java
+++ b/plc4j/integrations/apache-kafka/src/main/java/org/apache/plc4x/kafka/Plc4xSourceTask.java
@@ -28,13 +28,10 @@
import org.apache.plc4x.java.DefaultPlcDriverManager;
import org.apache.plc4x.java.api.PlcConnectionManager;
import org.apache.plc4x.java.api.value.PlcValue;
-import org.apache.plc4x.java.scraper.config.triggeredscraper.JobConfigurationTriggeredImplBuilder;
-import org.apache.plc4x.java.scraper.config.triggeredscraper.ScraperConfigurationTriggeredImpl;
-import org.apache.plc4x.java.scraper.config.triggeredscraper.ScraperConfigurationTriggeredImplBuilder;
-import org.apache.plc4x.java.scraper.exception.ScraperException;
-import org.apache.plc4x.java.scraper.triggeredscraper.TriggeredScraperImpl;
-import org.apache.plc4x.java.scraper.triggeredscraper.triggerhandler.collector.TriggerCollector;
-import org.apache.plc4x.java.scraper.triggeredscraper.triggerhandler.collector.TriggerCollectorImpl;
+import org.apache.plc4x.java.api.messages.PlcReadResponse;
+import org.apache.plc4x.java.tools.eventpump.EventPump;
+import org.apache.plc4x.java.tools.eventpump.TagBatch;
+import org.apache.plc4x.java.tools.eventpump.triggers.TimerTrigger;
import org.apache.plc4x.java.utils.cache.CachedPlcConnectionManager;
import org.apache.plc4x.kafka.config.Constants;
import org.apache.plc4x.kafka.util.VersionUtil;
@@ -93,10 +90,10 @@
.field(Constants.JOB_NAME_FIELD, Schema.STRING_SCHEMA)
.build();
- // Internal buffer into which all incoming scraper responses are written to.
+ // Internal buffer into which all incoming responses are written to.
private ArrayBlockingQueue<SourceRecord> buffer;
private Integer pollReturnInterval;
- private TriggeredScraperImpl scraper;
+ private EventPump eventPump;
private final SecureRandom random = new SecureRandom();
@Override
@@ -112,12 +109,14 @@
pollReturnInterval = config.getInt(Constants.KAFKA_POLL_RETURN_CONFIG);
Integer bufferSize = config.getInt(Constants.BUFFER_SIZE_CONFIG);
- Map<String, String> topics = new HashMap<>();
// Create a buffer with a capacity of BUFFER_SIZE_CONFIG elements which schedules access in a fair way.
buffer = new ArrayBlockingQueue<>(bufferSize, true);
- ScraperConfigurationTriggeredImplBuilder builder = new ScraperConfigurationTriggeredImplBuilder();
- builder.addSource(connectionName, plc4xConnectionString);
+ PlcConnectionManager connectionManager = CachedPlcConnectionManager.getBuilder()
+ .withConnectionManager(new DefaultPlcDriverManager())
+ .build();
+
+ eventPump = new EventPump();
List<String> jobConfigs = config.getList(Constants.QUERIES_CONFIG);
for (String jobConfig : jobConfigs) {
@@ -129,11 +128,11 @@
continue;
}
- String jobName = jobConfigSegments[0];
- String topic = jobConfigSegments[1];
- Integer rate = Integer.valueOf(jobConfigSegments[2]);
- JobConfigurationTriggeredImplBuilder jobBuilder = builder.job(
- jobName, String.format("(SCHEDULED,%s)", rate)).source(connectionName);
+ final String jobName = jobConfigSegments[0];
+ final String topic = jobConfigSegments[1];
+ final int rate = Integer.parseInt(jobConfigSegments[2]);
+
+ Map<String, String> tags = new LinkedHashMap<>();
for (int i = 3; i < jobConfigSegments.length; i++) {
String[] tagSegments = jobConfigSegments[i].split("#");
if (tagSegments.length != 2) {
@@ -142,108 +141,131 @@
jobName, jobConfigSegments[i]);
continue;
}
- String tagAlias = tagSegments[0];
- String tagAddress = tagSegments[1];
- jobBuilder.tag(tagAlias, tagAddress);
- topics.put(jobName, topic);
+ tags.put(tagSegments[0], tagSegments[1]);
}
- jobBuilder.build();
+ if (tags.isEmpty()) {
+ log.warn("Job configuration '{}' doesn't contain any valid tags ... skipping it.", jobName);
+ continue;
+ }
+
+ // One batch per job: all tags of a job are read together, at the job's rate.
+ TagBatch batch = TagBatch.builder()
+ .withBatchId(jobName)
+ .withConnectionManager(connectionManager)
+ .withConnectionString(plc4xConnectionString)
+ .addTagAddresses(tags)
+ .withTrigger(new TimerTrigger(rate, TimeUnit.MILLISECONDS))
+ .withListener(new TagBatch.TagBatchListener() {
+ @Override
+ public void onTagsFetched(TagBatch tagBatch, PlcReadResponse response) {
+ handleResponse(tagBatch.getBatchId(), connectionName, topic, response);
+ }
+
+ @Override
+ public void onError(TagBatch tagBatch, Throwable error) {
+ log.error("Error reading tags for job '{}': {}", tagBatch.getBatchId(), error.getMessage());
+ }
+
+ @Override
+ public void onFetchSkipped(TagBatch tagBatch, long lastFetchDurationMs, long consecutiveSkips) {
+ log.warn("Job '{}' is configured to be read every {}ms, but the last read took {}ms " +
+ "({} consecutive reads skipped).",
+ tagBatch.getBatchId(), rate, lastFetchDurationMs, consecutiveSkips);
+ }
+ })
+ .build();
+ eventPump.addBatch(batch);
}
- ScraperConfigurationTriggeredImpl scraperConfig = builder.build();
+ eventPump.startAll();
+ }
+ /**
+ * Turns one response into a Kafka {@link SourceRecord} and adds it to the buffer that
+ * {@link #poll()} drains.
+ */
+ private void handleResponse(String jobName, String sourceName, String topic, PlcReadResponse response) {
+ Map<String, Object> results = response.getTagNames().stream()
+ .collect(HashMap::new, (map, name) -> map.put(name, response.getObject(name)), HashMap::putAll);
try {
- PlcConnectionManager connectionManager = CachedPlcConnectionManager.getBuilder()
- .withConnectionManager(new DefaultPlcDriverManager())
+ Long timestamp = System.currentTimeMillis();
+
+ Map<String, String> sourcePartition = new HashMap<>();
+ sourcePartition.put("sourceName", sourceName);
+ sourcePartition.put("jobName", jobName);
+
+ Map<String, Long> sourceOffset = Collections.singletonMap("offset", timestamp);
+
+ // Prepare the key structure.
+ Struct key = new Struct(KEY_SCHEMA)
+ .put(Constants.SOURCE_NAME_FIELD, sourceName)
+ .put(Constants.JOB_NAME_FIELD, jobName);
+
+ // Build the Schema for the result struct.
+ SchemaBuilder tagSchemaBuilder = SchemaBuilder.struct()
+ .name("org.apache.plc4x.kafka.schema.Tag");
+
+
+ for (Map.Entry<String, Object> result : results.entrySet()) {
+ // Get tag-name and -value from the results.
+ String tagName = result.getKey();
+ Object tagValue = result.getValue();
+
+ // Get the schema for the given value type.
+ Schema valueSchema = getSchema(tagValue);
+
+ // Add the schema description for the current tag.
+ tagSchemaBuilder.field(tagName, valueSchema);
+ }
+ Schema tagSchema = tagSchemaBuilder.build();
+
+ Schema recordSchema = SchemaBuilder.struct()
+ .name("org.apache.plc4x.kafka.schema.JobResult")
+ .doc("PLC Job result. This contains all of the received PLCValues as well as a received timestamp")
+ .field(Constants.TAGS_CONFIG, tagSchema)
+ .field(Constants.TIMESTAMP_CONFIG, Schema.INT64_SCHEMA)
+ .field(Constants.EXPIRES_CONFIG, Schema.OPTIONAL_INT64_SCHEMA)
.build();
- TriggerCollector triggerCollector = new TriggerCollectorImpl(connectionManager);
- scraper = new TriggeredScraperImpl(scraperConfig, connectionManager, (jobName, sourceName, results) -> {
- try {
- Long timestamp = System.currentTimeMillis();
- Map<String, String> sourcePartition = new HashMap<>();
- sourcePartition.put("sourceName", sourceName);
- sourcePartition.put("jobName", jobName);
+ // Build the struct itself.
+ Struct tagStruct = new Struct(tagSchema);
+ for (Map.Entry<String, Object> result : results.entrySet()) {
+ // Get tag-name and -value from the results.
+ String tagName = result.getKey();
+ Object tagValue = result.getValue();
- Map<String, Long> sourceOffset = Collections.singletonMap("offset", timestamp);
-
- String topic = topics.get(jobName);
-
- // Prepare the key structure.
- Struct key = new Struct(KEY_SCHEMA)
- .put(Constants.SOURCE_NAME_FIELD, sourceName)
- .put(Constants.JOB_NAME_FIELD, jobName);
-
- // Build the Schema for the result struct.
- SchemaBuilder tagSchemaBuilder = SchemaBuilder.struct()
- .name("org.apache.plc4x.kafka.schema.Tag");
-
-
- for (Map.Entry<String, Object> result : results.entrySet()) {
- // Get tag-name and -value from the results.
- String tagName = result.getKey();
- Object tagValue = result.getValue();
-
- // Get the schema for the given value type.
- Schema valueSchema = getSchema(tagValue);
-
- // Add the schema description for the current tag.
- tagSchemaBuilder.field(tagName, valueSchema);
- }
- Schema tagSchema = tagSchemaBuilder.build();
-
- Schema recordSchema = SchemaBuilder.struct()
- .name("org.apache.plc4x.kafka.schema.JobResult")
- .doc("PLC Job result. This contains all of the received PLCValues as well as a received timestamp")
- .field(Constants.TAGS_CONFIG, tagSchema)
- .field(Constants.TIMESTAMP_CONFIG, Schema.INT64_SCHEMA)
- .field(Constants.EXPIRES_CONFIG, Schema.OPTIONAL_INT64_SCHEMA)
- .build();
-
- // Build the struct itself.
- Struct tagStruct = new Struct(tagSchema);
- for (Map.Entry<String, Object> result : results.entrySet()) {
- // Get tag-name and -value from the results.
- String tagName = result.getKey();
- Object tagValue = result.getValue();
-
- if (tagSchema.field(tagName).schema().type() == Schema.Type.ARRAY) {
- tagStruct.put(tagName, ((List) tagValue).stream().map(p -> ((PlcValue) p).getObject()).collect(Collectors.toList()));
- } else {
- tagStruct.put(tagName, tagValue);
- }
- }
-
- Struct recordStruct = new Struct(recordSchema)
- .put(Constants.TAGS_CONFIG, tagStruct)
- .put(Constants.TIMESTAMP_CONFIG, timestamp);
-
- // Prepare the source-record element.
- SourceRecord sourceRecord = new SourceRecord(
- sourcePartition, sourceOffset,
- topic,
- KEY_SCHEMA, key,
- recordSchema, recordStruct
- );
-
- // Add the new source-record to the buffer.
- buffer.add(sourceRecord);
- } catch (Exception e) {
- log.error("Error while parsing returned values", e);
+ if (tagSchema.field(tagName).schema().type() == Schema.Type.ARRAY) {
+ tagStruct.put(tagName, ((List) tagValue).stream().map(p -> ((PlcValue) p).getObject()).collect(Collectors.toList()));
+ } else {
+ tagStruct.put(tagName, tagValue);
}
- }, triggerCollector);
- scraper.start();
- triggerCollector.start();
- } catch (ScraperException e) {
- log.error("Error starting the scraper", e);
+ }
+ Struct recordStruct = new Struct(recordSchema)
+ .put(Constants.TAGS_CONFIG, tagStruct)
+ .put(Constants.TIMESTAMP_CONFIG, timestamp);
+
+ // Prepare the source-record element.
+ SourceRecord sourceRecord = new SourceRecord(
+ sourcePartition, sourceOffset,
+ topic,
+ KEY_SCHEMA, key,
+ recordSchema, recordStruct
+ );
+
+ // Add the new source-record to the buffer.
+ buffer.add(sourceRecord);
+ } catch (Exception e) {
+ log.error("Error while parsing returned values", e);
}
}
@Override
public void stop() {
synchronized (this) {
- scraper.stop();
+ if (eventPump != null) {
+ eventPump.close();
+ }
notifyAll(); // wake up thread waiting in awaitFetch
}
}
diff --git a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SinkTaskTest.java b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SinkTaskTest.java
index 27b20a3..e1c80d7 100644
--- a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SinkTaskTest.java
+++ b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SinkTaskTest.java
@@ -79,7 +79,7 @@
assertEquals("machineSinkA", config.get(0).get(Constants.TOPIC_CONFIG));
assertEquals("5", config.get(0).get(Constants.RETRIES_CONFIG));
assertEquals("5000", config.get(0).get(Constants.TIMEOUT_CONFIG));
- assertEquals("running#RANDOM/Temporary:Boolean|conveyorEntry#RANDOM/Temporary:Boolean|load#RANDOM/Temporary:Boolean|unload#RANDOM/Temporary:Boolean|transferLeft#RANDOM/Temporary:Boolean|transferRight#RANDOM/Temporary:Boolean|conveyorLeft#RANDOM/Temporary:Boolean|conveyorRight#RANDOM/Temporary:Boolean|numLargeBoxes#STATE/Temporary:Integer|numSmallBoxes#RANDOM/Temporary:Integer",
+ assertEquals("running#RANDOM/Temporary:BOOL|conveyorEntry#RANDOM/Temporary:BOOL|load#RANDOM/Temporary:BOOL|unload#RANDOM/Temporary:BOOL|transferLeft#RANDOM/Temporary:BOOL|transferRight#RANDOM/Temporary:BOOL|conveyorLeft#RANDOM/Temporary:BOOL|conveyorRight#RANDOM/Temporary:BOOL|numLargeBoxes#STATE/Temporary:Integer|numSmallBoxes#RANDOM/Temporary:DINT",
config.get(0).get(Constants.QUERIES_CONFIG));
}
diff --git a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SourceTaskTest.java b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SourceTaskTest.java
index 2fdfe23..dae0111 100644
--- a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SourceTaskTest.java
+++ b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/config/SourceTaskTest.java
@@ -20,6 +20,7 @@
import static org.junit.jupiter.api.Assertions.*;
+import org.apache.kafka.connect.source.SourceRecord;
import org.apache.plc4x.kafka.Plc4xSourceConnector;
import org.apache.plc4x.kafka.Plc4xSourceTask;
@@ -69,7 +70,7 @@
assertEquals("simulated://127.0.0.1", config.get(0).get(Constants.CONNECTION_STRING_CONFIG));
assertEquals("1000", config.get(0).get(Constants.BUFFER_SIZE_CONFIG));
assertEquals("5000", config.get(0).get(Constants.KAFKA_POLL_RETURN_CONFIG));
- assertEquals("simulateddashboard|machineData|1000|running#RANDOM/Temporary:Boolean|conveyorEntry#RANDOM/Temporary:Boolean|load#RANDOM/Temporary:Boolean|unload#RANDOM/Temporary:Boolean|transferLeft#RANDOM/Temporary:Boolean|transferRight#RANDOM/Temporary:Boolean|conveyorLeft#RANDOM/Temporary:Boolean|conveyorRight#RANDOM/Temporary:Boolean|numLargeBoxes#RANDOM/Temporary:Integer|numSmallBoxes#RANDOM/Temporary:Integer[2],simulatedheartbeat|simulatedheartbeat|500|active#RANDOM/Temporary:Integer", config.get(0).get(Constants.QUERIES_CONFIG));
+ assertEquals("simulateddashboard|machineData|1000|running#RANDOM/Temporary:BOOL|conveyorEntry#RANDOM/Temporary:BOOL|load#RANDOM/Temporary:BOOL|unload#RANDOM/Temporary:BOOL|transferLeft#RANDOM/Temporary:BOOL|transferRight#RANDOM/Temporary:BOOL|conveyorLeft#RANDOM/Temporary:BOOL|conveyorRight#RANDOM/Temporary:BOOL|numLargeBoxes#RANDOM/Temporary:DINT|numSmallBoxes#RANDOM/Temporary:DINT[2],simulatedheartbeat|simulatedheartbeat|500|active#RANDOM/Temporary:DINT", config.get(0).get(Constants.QUERIES_CONFIG));
}
@Test
@@ -85,8 +86,18 @@
sourceTask.start(taskConfig);
}
Thread.sleep(5000);
- for (Plc4xSourceTask sourceTask : sourceList) {
- assertNotNull(sourceTask.poll());
+ try {
+ for (Plc4xSourceTask sourceTask : sourceList) {
+ List<SourceRecord> records = sourceTask.poll();
+ assertNotNull(records);
+ // The jobs are configured to be read every 1000ms/500ms, so after 5s there
+ // has to be data in the buffer - an empty result means nothing was ever read.
+ assertFalse(records.isEmpty(), "expected the source task to have collected records");
+ }
+ } finally {
+ for (Plc4xSourceTask sourceTask : sourceList) {
+ sourceTask.stop();
+ }
}
}
diff --git a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/sink_task_no_error.properties b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/sink_task_no_error.properties
index 68b31f3..ea26c4f 100644
--- a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/sink_task_no_error.properties
+++ b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/sink_task_no_error.properties
@@ -26,24 +26,24 @@
sinks.machineA.retries=5
sinks.machineA.tags=running,conveyorEntry,load,unload,transferLeft,transferRight,conveyorLeft,conveyorRight,numLargeBoxes,numSmallBoxes
-sinks.machineA.tags.running=RANDOM/Temporary:Boolean
-sinks.machineA.tags.conveyorEntry=RANDOM/Temporary:Boolean
-sinks.machineA.tags.load=RANDOM/Temporary:Boolean
-sinks.machineA.tags.unload=RANDOM/Temporary:Boolean
-sinks.machineA.tags.transferLeft=RANDOM/Temporary:Boolean
-sinks.machineA.tags.transferRight=RANDOM/Temporary:Boolean
-sinks.machineA.tags.conveyorLeft=RANDOM/Temporary:Boolean
-sinks.machineA.tags.conveyorRight=RANDOM/Temporary:Boolean
+sinks.machineA.tags.running=RANDOM/Temporary:BOOL
+sinks.machineA.tags.conveyorEntry=RANDOM/Temporary:BOOL
+sinks.machineA.tags.load=RANDOM/Temporary:BOOL
+sinks.machineA.tags.unload=RANDOM/Temporary:BOOL
+sinks.machineA.tags.transferLeft=RANDOM/Temporary:BOOL
+sinks.machineA.tags.transferRight=RANDOM/Temporary:BOOL
+sinks.machineA.tags.conveyorLeft=RANDOM/Temporary:BOOL
+sinks.machineA.tags.conveyorRight=RANDOM/Temporary:BOOL
sinks.machineA.tags.numLargeBoxes=STATE/Temporary:Integer
-sinks.machineA.tags.numSmallBoxes=RANDOM/Temporary:Integer
+sinks.machineA.tags.numSmallBoxes=RANDOM/Temporary:DINT
sinks.machineB.connectionString=simulated://127.0.0.1
sinks.machineB.topic=machineSinkB
sinks.machineB.timeout=5000
sinks.machineB.retries=5
sinks.machineB.tags=running,conveyorEntry
-sinks.machineB.tags.running=RANDOM/Temporary:Boolean
-sinks.machineB.tags.conveyorEntry=RANDOM/Temporary:Boolean
+sinks.machineB.tags.running=RANDOM/Temporary:BOOL
+sinks.machineB.tags.conveyorEntry=RANDOM/Temporary:BOOL
bootstrap.servers=127.0.0.1:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
diff --git a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/source_task_no_error.properties b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/source_task_no_error.properties
index adf250e..a8f0a79 100644
--- a/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/source_task_no_error.properties
+++ b/plc4j/integrations/apache-kafka/src/test/java/org/apache/plc4x/kafka/properties/source_task_no_error.properties
@@ -31,20 +31,20 @@
jobs.simulateddashboard.interval=1000
jobs.simulateddashboard.tags=running,conveyorEntry,load,unload,transferLeft,transferRight,conveyorLeft,conveyorRight,numLargeBoxes,numSmallBoxes
-jobs.simulateddashboard.tags.running=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.conveyorEntry=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.load=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.unload=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.transferLeft=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.transferRight=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.conveyorLeft=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.conveyorRight=RANDOM/Temporary:Boolean
-jobs.simulateddashboard.tags.numLargeBoxes=RANDOM/Temporary:Integer
-jobs.simulateddashboard.tags.numSmallBoxes=RANDOM/Temporary:Integer[2]
+jobs.simulateddashboard.tags.running=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.conveyorEntry=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.load=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.unload=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.transferLeft=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.transferRight=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.conveyorLeft=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.conveyorRight=RANDOM/Temporary:BOOL
+jobs.simulateddashboard.tags.numLargeBoxes=RANDOM/Temporary:DINT
+jobs.simulateddashboard.tags.numSmallBoxes=RANDOM/Temporary:DINT[2]
jobs.simulatedheartbeat.interval=500
jobs.simulatedheartbeat.tags=active
-jobs.simulatedheartbeat.tags.active=RANDOM/Temporary:Integer
+jobs.simulatedheartbeat.tags.active=RANDOM/Temporary:DINT
bootstrap.servers=localhost:9092
#key.converter=org.apache.kafka.connect.json.JsonConverter