This is an automated email from the ASF dual-hosted git repository.
dominikriemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new af2fe9e1dd fix: Reduce logging for time series storage and adapter
preprocessing (#4657)
af2fe9e1dd is described below
commit af2fe9e1ddf64af114c8965cd11f44cd6b33e096
Author: Dominik Riemer <[email protected]>
AuthorDate: Wed Jul 1 13:30:59 2026 +0200
fix: Reduce logging for time series storage and adapter preprocessing
(#4657)
---
.../shared/AdapterPipelineGeneratorBase.java | 1 +
.../streampipes/connect/shared/DatatypeUtils.java | 124 +++++++++++++++++----
.../value/DatatypeTransformationRule.java | 10 +-
.../preprocessing/transform/DatatypeUtilsTest.java | 71 ++++++++++--
.../dataexplorer/TimeSeriesStorage.java | 23 +++-
.../connect/adapter/parser/CsvParser.java | 4 +-
6 files changed, 196 insertions(+), 37 deletions(-)
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/AdapterPipelineGeneratorBase.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/AdapterPipelineGeneratorBase.java
index 9c89c7ce8b..f770792d01 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/AdapterPipelineGeneratorBase.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/AdapterPipelineGeneratorBase.java
@@ -60,6 +60,7 @@ public class AdapterPipelineGeneratorBase {
.filter(ep -> ep.getAdditionalMetadata()
.containsKey("originType"))
.map(ep -> new DatatypeTransformationRule(
+ adapterDescription.getName(),
ep.getRuntimeName(),
((EventPropertyPrimitive) ep).getRuntimeType()
))
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/DatatypeUtils.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/DatatypeUtils.java
index 1f95d29eef..29b0a077e4 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/DatatypeUtils.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/DatatypeUtils.java
@@ -24,6 +24,8 @@ import org.apache.commons.lang3.math.NumberUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.concurrent.atomic.AtomicBoolean;
+
public class DatatypeUtils {
private static final Logger LOG =
LoggerFactory.getLogger(DatatypeUtils.class);
@@ -35,40 +37,116 @@ public class DatatypeUtils {
* If the conversion is not possible due to a format mismatch, the original
value is returned.
* A number format exception during conversion is logged as an error.
*
+ * @param adapterName The adapter whose event value should be converted.
* @param value The value to be converted. It can be of any type.
* @param targetDatatypeXsd The target XSD datatype as a string. Supported
types are XSD.STRING,
* XSD.DOUBLE, XSD.FLOAT, XSD.BOOLEAN, XSD.INTEGER,
and XSD.LONG.
* @return The converted value as an Object. If conversion fails, the
original value is returned.
- * @throws NumberFormatException if the string does not contain a parsable
number for numeric conversions.
*/
- public static Object convertValue(Object value,
+ public static Object convertValue(String adapterName,
+ Object value,
String targetDatatypeXsd) {
- var stringValue = String.valueOf(value);
+ return convertValue(adapterName, value, targetDatatypeXsd, new
AtomicBoolean(false));
+ }
+
+ public static Object convertValue(String adapterName,
+ Object value,
+ String targetDatatypeXsd,
+ AtomicBoolean loggedConversionError) {
+ if (value == null) {
+ return null;
+ }
+
if (XSD.STRING.toString().equals(targetDatatypeXsd)) {
- return stringValue;
+ return String.valueOf(value);
+ }
+
+ if (value instanceof Number number &&
isNumericDatatype(targetDatatypeXsd)) {
+ return convertNumber(number, targetDatatypeXsd);
+ }
+
+ if (value instanceof Boolean booleanValue &&
XSD.BOOLEAN.toString().equals(targetDatatypeXsd)) {
+ return booleanValue;
+ }
+
+ if (!isSupportedDatatype(targetDatatypeXsd)) {
+ return value;
+ }
+
+ try {
+ return convertString(String.valueOf(value), targetDatatypeXsd);
+ } catch (NumberFormatException e) {
+ logConversionError(adapterName, value, targetDatatypeXsd,
loggedConversionError);
+ return value;
+ }
+ }
+
+ private static void logConversionError(String adapterName,
+ Object value,
+ String targetDatatypeXsd,
+ AtomicBoolean loggedConversionError) {
+ if (loggedConversionError.compareAndSet(false, true)) {
+ LOG.warn(
+ "Could not convert value '{}' to datatype '{}' for adapter '{}'.
Further occurrences are logged at debug "
+ + "level.",
+ value,
+ targetDatatypeXsd,
+ adapterName
+ );
} else {
- try {
- if (XSD.DOUBLE.toString().equals(targetDatatypeXsd)) {
- return Double.parseDouble(stringValue);
- } else if (XSD.FLOAT.toString().equals(targetDatatypeXsd)) {
- return Float.parseFloat(stringValue);
- } else if (XSD.BOOLEAN.toString().equals(targetDatatypeXsd)) {
- return Boolean.parseBoolean(stringValue);
- } else if (XSD.INTEGER.toString().equals(targetDatatypeXsd)) {
- return ((Double) Double.parseDouble(stringValue)).intValue();
- } else if (XSD.LONG.toString().equals(targetDatatypeXsd)) {
- var floatingNumber = Double.parseDouble(stringValue);
- return Long.parseLong(String.valueOf(Math.round(floatingNumber)));
- }
- } catch (NumberFormatException e) {
- LOG.error("Number format exception {}", value);
- return value;
- }
+ LOG.debug(
+ "Could not convert value '{}' to datatype '{}' for adapter '{}'",
+ value,
+ targetDatatypeXsd,
+ adapterName
+ );
+ }
+ }
+
+ private static Object convertNumber(Number value,
+ String targetDatatypeXsd) {
+ if (XSD.DOUBLE.toString().equals(targetDatatypeXsd)) {
+ return value.doubleValue();
+ } else if (XSD.FLOAT.toString().equals(targetDatatypeXsd)) {
+ return value.floatValue();
+ } else if (XSD.INTEGER.toString().equals(targetDatatypeXsd)) {
+ return value.intValue();
+ } else if (XSD.LONG.toString().equals(targetDatatypeXsd)) {
+ return Math.round(value.doubleValue());
+ }
+
+ return value;
+ }
+
+ private static Object convertString(String value,
+ String targetDatatypeXsd) {
+ if (XSD.DOUBLE.toString().equals(targetDatatypeXsd)) {
+ return Double.parseDouble(value);
+ } else if (XSD.FLOAT.toString().equals(targetDatatypeXsd)) {
+ return Float.parseFloat(value);
+ } else if (XSD.BOOLEAN.toString().equals(targetDatatypeXsd)) {
+ return Boolean.parseBoolean(value);
+ } else if (XSD.INTEGER.toString().equals(targetDatatypeXsd)) {
+ return ((Double) Double.parseDouble(value)).intValue();
+ } else if (XSD.LONG.toString().equals(targetDatatypeXsd)) {
+ var floatingNumber = Double.parseDouble(value);
+ return Long.parseLong(String.valueOf(Math.round(floatingNumber)));
}
return value;
}
+ private static boolean isSupportedDatatype(String targetDatatypeXsd) {
+ return isNumericDatatype(targetDatatypeXsd) ||
XSD.BOOLEAN.toString().equals(targetDatatypeXsd);
+ }
+
+ private static boolean isNumericDatatype(String targetDatatypeXsd) {
+ return XSD.DOUBLE.toString().equals(targetDatatypeXsd)
+ || XSD.FLOAT.toString().equals(targetDatatypeXsd)
+ || XSD.INTEGER.toString().equals(targetDatatypeXsd)
+ || XSD.LONG.toString().equals(targetDatatypeXsd);
+ }
+
public static String getXsdDatatype(String value,
boolean preferFloat) {
var clazz = getTypeClass(value, preferFloat);
@@ -90,6 +168,10 @@ public class DatatypeUtils {
public static Class<?> getTypeClass(String value,
boolean preferFloatingPointNumber) {
var targetClass = String.class;
+ if (value == null) {
+ return targetClass;
+ }
+
if (NumberUtils.isParsable(value)) {
Class<?> numberClass;
try {
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
index 59d0144f47..534dcc73f3 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
@@ -23,15 +23,21 @@ import org.apache.streampipes.connect.shared.DatatypeUtils;
import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
public class DatatypeTransformationRule implements TransformationRule {
private final String eventKey;
private String targetDatatypeXsd;
+ private final String adapterName;
+ private final AtomicBoolean loggedConversionError = new AtomicBoolean(false);
- public DatatypeTransformationRule(String eventKey,
+ public DatatypeTransformationRule(String adapterName,
+ String eventKey,
String targetDatatypeXsd) {
+
this.eventKey = eventKey;
+ this.adapterName = adapterName;
this.targetDatatypeXsd = targetDatatypeXsd;
}
@@ -45,6 +51,6 @@ public class DatatypeTransformationRule implements
TransformationRule {
}
public Object transformDatatype(Object value) {
- return DatatypeUtils.convertValue(value, targetDatatypeXsd);
+ return DatatypeUtils.convertValue(adapterName, value, targetDatatypeXsd,
loggedConversionError);
}
}
diff --git
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/DatatypeUtilsTest.java
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/DatatypeUtilsTest.java
index 7937d6875c..f9d62a2278 100644
---
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/DatatypeUtilsTest.java
+++
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/DatatypeUtilsTest.java
@@ -26,9 +26,13 @@ import org.junit.jupiter.api.Test;
import java.util.Locale;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
public class DatatypeUtilsTest {
+ private static final String TEST_ADAPTER_NAME = "test-adapter";
+
/**
* The following tests ensure that timestamps represented as strings are
correctly parsed.
* Often they are first parsed into floating point number before transformed
back to long.
@@ -38,88 +42,124 @@ public class DatatypeUtilsTest {
@Test
public void convertValue_StringToStringValue() {
var inputValue = "testString";
- var actualValue = DatatypeUtils.convertValue(inputValue,
XSD.STRING.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
inputValue, XSD.STRING.toString());
assertEquals(inputValue, actualValue);
}
+ @Test
+ public void convertValue_NullToStringValue() {
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, null,
XSD.STRING.toString());
+
+ assertNull(actualValue);
+ }
+
+ @Test
+ public void convertValue_NullToNumericValue() {
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, null,
XSD.DOUBLE.toString());
+
+ assertNull(actualValue);
+ }
+
@Test
public void convertValue_StringToDoubleValue() {
- var actualValue = DatatypeUtils.convertValue("1667904471000",
XSD.DOUBLE.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
"1667904471000", XSD.DOUBLE.toString());
assertEquals(1.667904471E12, actualValue);
}
+ @Test
+ public void convertValue_IntegerToDoubleValue() {
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, 123,
XSD.DOUBLE.toString());
+
+ assertEquals(123.0d, actualValue);
+ }
+
@Test
public void convertValue_StringToFloatValue() {
- var actualValue = DatatypeUtils.convertValue("123.45",
XSD.FLOAT.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, "123.45",
XSD.FLOAT.toString());
assertEquals(123.45f, actualValue);
}
@Test
public void convertValue_StringToInteger() {
- var actualValue = DatatypeUtils.convertValue("1623871500",
XSD.INTEGER.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
"1623871500", XSD.INTEGER.toString());
assertEquals(1623871500, actualValue);
}
@Test
public void convertValue_StringToIntegerValue() {
- var actualValue = DatatypeUtils.convertValue("123",
XSD.INTEGER.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, "123",
XSD.INTEGER.toString());
assertEquals(123, actualValue);
}
@Test
public void convertValue_StringToLongValue() {
- var actualValue = DatatypeUtils.convertValue("1623871500000",
XSD.LONG.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
"1623871500000", XSD.LONG.toString());
assertEquals(1623871500000L, actualValue);
}
@Test
public void convertValue_StringToBooleanTrueValue() {
- var actualValue = DatatypeUtils.convertValue("true",
XSD.BOOLEAN.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, "true",
XSD.BOOLEAN.toString());
+
+ assertEquals(true, actualValue);
+ }
+
+ @Test
+ public void convertValue_BooleanToBooleanValue() {
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, true,
XSD.BOOLEAN.toString());
assertEquals(true, actualValue);
}
@Test
public void convertValue_StringToBooleanFalseValue() {
- var actualValue = DatatypeUtils.convertValue("false",
XSD.BOOLEAN.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, "false",
XSD.BOOLEAN.toString());
assertEquals(false, actualValue);
}
@Test
public void convertValue_FloatToIntegerValue_Rounding() {
- var actualValue = DatatypeUtils.convertValue(123.45f,
XSD.INTEGER.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME, 123.45f,
XSD.INTEGER.toString());
assertEquals(123, actualValue);
}
@Test
public void convertValue_DoubleToLongValue_Rounding1() {
- var actualValue = DatatypeUtils.convertValue(1234567890.12345,
XSD.LONG.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
1234567890.12345, XSD.LONG.toString());
assertEquals(1234567890L, actualValue);
}
@Test
public void convertValue_DoubleToLongValue() {
- var actualValue = DatatypeUtils.convertValue(1.667904471E12,
XSD.LONG.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
1.667904471E12, XSD.LONG.toString());
assertEquals(1667904471000L, actualValue);
}
@Test
public void convertValue_DoubleToLongValue_Rounding() {
- var actualValue = DatatypeUtils.convertValue(1234567890.12345,
XSD.LONG.toString());
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
1234567890.12345, XSD.LONG.toString());
assertEquals(1234567890L, actualValue);
}
+ @Test
+ public void convertValue_UnsupportedDatatypeReturnsOriginalValue() {
+ var inputValue = new Object();
+ var actualValue = DatatypeUtils.convertValue(TEST_ADAPTER_NAME,
inputValue, "unsupported");
+
+ assertSame(inputValue, actualValue);
+ }
+
String booleanInputValue = "true";
@@ -209,4 +249,11 @@ public class DatatypeUtilsTest {
assertEquals(String.class, result);
}
+ @Test
+ public void getTypeClass_Null() {
+ var result = DatatypeUtils.getTypeClass(null, false);
+
+ assertEquals(String.class, result);
+ }
+
}
diff --git
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
index ceaa6b9dc0..850760df82 100644
---
a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
+++
b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/TimeSeriesStorage.java
@@ -31,6 +31,7 @@ import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
public abstract class TimeSeriesStorage implements ITimeSeriesStorage {
@@ -40,6 +41,7 @@ public abstract class TimeSeriesStorage implements
ITimeSeriesStorage {
protected final DataLakeMeasure measure;
protected final List<EventProperty> allEventProperties;
protected final Map<String, String> sanitizedRuntimeNames = new HashMap<>();
+ private final AtomicBoolean warnedNullFields = new AtomicBoolean(false);
public TimeSeriesStorage(DataLakeMeasure measure) {
this.measure = measure;
@@ -115,7 +117,26 @@ public abstract class TimeSeriesStorage implements
ITimeSeriesStorage {
.collect(Collectors.toList());
if (!nullFields.isEmpty()) {
- LOG.warn("Ignored {} fields which had a value 'null': {}",
nullFields.size(), String.join(", ", nullFields));
+ logNullFields(nullFields);
+ }
+ }
+
+ private void logNullFields(List<String> nullFields) {
+ var nullFieldsLogMessage = String.join(", ", nullFields);
+
+ if (warnedNullFields.compareAndSet(false, true)) {
+ LOG.warn(
+ "Ignored {} fields which had a value 'null': {}. Further occurrences
for this measure are logged at "
+ + "debug level.",
+ nullFields.size(),
+ nullFieldsLogMessage
+ );
+ } else {
+ LOG.debug(
+ "Ignored {} fields which had a value 'null': {}",
+ nullFields.size(),
+ nullFieldsLogMessage
+ );
}
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/parser/CsvParser.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/parser/CsvParser.java
index b55e2a2f27..3a4da78e94 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/parser/CsvParser.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/parser/CsvParser.java
@@ -50,6 +50,7 @@ import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.IntStream;
public class CsvParser implements IParser {
@@ -67,6 +68,7 @@ public class CsvParser implements IParser {
private boolean header;
private char delimiter;
+ private final AtomicBoolean loggedConversionError = new AtomicBoolean(false);
public CsvParser() {
@@ -163,7 +165,7 @@ public class CsvParser implements IParser {
var event = new HashMap<String, Object>();
for (int i = 0; i < header.length; i++) {
var runtimeType = DatatypeUtils.getXsdDatatype(values[i], preferFloat);
- var convertedValue = DatatypeUtils.convertValue(values[i], runtimeType);
+ var convertedValue = DatatypeUtils.convertValue(ID, values[i],
runtimeType, loggedConversionError);
event.put(header[i], convertedValue);
}