This is an automated email from the ASF dual-hosted git repository.
zehnder pushed a commit to branch
2964-timestamp-conversion-is-broken-in-adapters
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to
refs/heads/2964-timestamp-conversion-is-broken-in-adapters by this push:
new 644a319132 refactor(#2964): Refactor method onAdapterStarted in
FileReplayAdapter
644a319132 is described below
commit 644a31913259dc44b4c15370ae1275ff07489726
Author: Philipp Zehnder <[email protected]>
AuthorDate: Mon Jul 1 15:20:49 2024 +0200
refactor(#2964): Refactor method onAdapterStarted in FileReplayAdapter
---
.../iiot/protocol/stream/FileReplayAdapter.java | 55 ++++++++++++++--------
1 file changed, 35 insertions(+), 20 deletions(-)
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
index 6a8b14650f..cdc3fb371d 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
@@ -135,8 +135,39 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
IAdapterRuntimeContext adapterRuntimeContext
) throws AdapterException {
- // extract user input
+ boolean replayOnce =
extractUserInputsAndReturnValueOfReplayOnce(extractor);
+
+ determineTimestampRuntimeName(extractor);
+
+ determineSourceTimestampField(extractor);
+
+ startAdapterReplayThread(extractor, collector, adapterRuntimeContext,
replayOnce);
+ }
+
+ private void startAdapterReplayThread(
+ IAdapterParameterExtractor extractor,
+ IEventCollector collector,
+ IAdapterRuntimeContext adapterRuntimeContext,
+ boolean replayOnce
+ ) {
executor = Executors.newScheduledThreadPool(1);
+ if (replayOnce) {
+ executor.schedule(
+ () -> getFileFromEndpointAndParseFile(extractor, collector,
adapterRuntimeContext),
+ 0,
+ TimeUnit.SECONDS
+ );
+ } else {
+ executor.scheduleAtFixedRate(
+ () -> getFileFromEndpointAndParseFile(extractor, collector,
adapterRuntimeContext),
+ 0,
+ 1,
+ TimeUnit.SECONDS
+ );
+ }
+ }
+
+ private boolean
extractUserInputsAndReturnValueOfReplayOnce(IAdapterParameterExtractor
extractor) {
boolean replayOnce = extractor
.getStaticPropertyExtractor()
.selectedSingleValue(REPLAY_ONCE, String.class)
@@ -155,8 +186,10 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
.singleValueParameter(SPEED_UP, Float.class);
default -> 1.0f;
};
+ return replayOnce;
+ }
- // get timestamp field
+ private void determineTimestampRuntimeName(IAdapterParameterExtractor
extractor) throws AdapterException {
var timestampField = extractor
.getAdapterDescription()
.getEventSchema()
@@ -173,24 +206,6 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
timestampRuntimeName = timestampField.get()
.getRuntimeName();
}
-
- determineSourceTimestampField(extractor);
-
- // start replay adapter
- if (replayOnce) {
- executor.schedule(
- () -> getFileFromEndpointAndParseFile(extractor, collector,
adapterRuntimeContext),
- 0,
- TimeUnit.SECONDS
- );
- } else {
- executor.scheduleAtFixedRate(
- () -> getFileFromEndpointAndParseFile(extractor, collector,
adapterRuntimeContext),
- 0,
- 1,
- TimeUnit.SECONDS
- );
- }
}
/**