This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch reduce-logging-extensions in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 2901847107170df68a3bd9f50779bf6c8431cead Author: Dominik Riemer <[email protected]> AuthorDate: Mon Jun 29 21:28:22 2026 +0200 fix: Reduce logging for time series storage and adapter preprocessing --- .../shared/AdapterPipelineGeneratorBase.java | 1 + .../streampipes/connect/shared/DatatypeUtils.java | 126 +++++++++++++++++---- .../value/DatatypeTransformationRule.java | 11 +- .../preprocessing/transform/DatatypeUtilsTest.java | 71 ++++++++++-- .../dataexplorer/TimeSeriesStorage.java | 25 +++- .../connect/adapter/parser/CsvParser.java | 5 +- 6 files changed, 202 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..18c761b1f3 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,9 @@ import org.apache.commons.lang3.math.NumberUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + public class DatatypeUtils { private static final Logger LOG = LoggerFactory.getLogger(DatatypeUtils.class); @@ -35,40 +38,117 @@ 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, ConcurrentHashMap.newKeySet()); + } + + public static Object convertValue(String adapterName, + Object value, + String targetDatatypeXsd, + Set<String> loggedConversionErrors) { + 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, loggedConversionErrors); + return value; + } + } + + private static void logConversionError(String adapterName, + Object value, + String targetDatatypeXsd, + Set<String> loggedConversionErrors) { + var logKey = "%s:%s:%s".formatted(adapterName, targetDatatypeXsd, value); + if (loggedConversionErrors.add(logKey)) { + 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 +170,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..3634a0f430 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,22 @@ import org.apache.streampipes.connect.shared.DatatypeUtils; import org.apache.streampipes.extensions.api.connect.TransformationRule; import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; public class DatatypeTransformationRule implements TransformationRule { private final String eventKey; private String targetDatatypeXsd; + private final String adapterName; + private final Set<String> loggedConversionErrors = ConcurrentHashMap.newKeySet(); - public DatatypeTransformationRule(String eventKey, + public DatatypeTransformationRule(String adapterName, + String eventKey, String targetDatatypeXsd) { + this.eventKey = eventKey; + this.adapterName = adapterName; this.targetDatatypeXsd = targetDatatypeXsd; } @@ -45,6 +52,6 @@ public class DatatypeTransformationRule implements TransformationRule { } public Object transformDatatype(Object value) { - return DatatypeUtils.convertValue(value, targetDatatypeXsd); + return DatatypeUtils.convertValue(adapterName, value, targetDatatypeXsd, loggedConversionErrors); } } 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..ef87811cc4 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,8 @@ import org.slf4j.LoggerFactory; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; public abstract class TimeSeriesStorage implements ITimeSeriesStorage { @@ -40,6 +42,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 Set<String> warnedNullFields = ConcurrentHashMap.newKeySet(); public TimeSeriesStorage(DataLakeMeasure measure) { this.measure = measure; @@ -115,7 +118,27 @@ 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); + var logKey = "%s:null-fields".formatted(measure.getMeasureName()); + + if (warnedNullFields.add(logKey)) { + 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..190721a004 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,8 @@ import java.util.Arrays; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.stream.IntStream; public class CsvParser implements IParser { @@ -67,6 +69,7 @@ public class CsvParser implements IParser { private boolean header; private char delimiter; + private final Set<String> loggedConversionErrors = ConcurrentHashMap.newKeySet(); public CsvParser() { @@ -163,7 +166,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, loggedConversionErrors); event.put(header[i], convertedValue); }
