This is an automated email from the ASF dual-hosted git repository.
mattcasters pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new 2e67685c0c Issue #8600 : Save Avro Record values in an execution
information location (#8607)
2e67685c0c is described below
commit 2e67685c0cb762989c7ed96b80964039e0e45329
Author: Matt Casters <[email protected]>
AuthorDate: Mon Sep 28 12:26:42 2026 +0200
Issue #8600 : Save Avro Record values in an execution information location
(#8607)
Store the writer schema and a length-prefixed Avro datum with each value so
a
field that has no schema on its metadata can be read back. An execution that
was already saved with the old layout still opens.
---
.../hop/core/row/value/ValueMetaAvroRecord.java | 149 ++++++++++++++++++---
.../core/row/value/ValueMetaAvroRecordTest.java | 47 +++++++
.../org/apache/hop/execution/ExecutionData.java | 58 +++++---
.../dataprof/BasicDataProfilingDataSampler.java | 2 +-
.../apache/hop/execution/ExecutionDataTest.java | 117 ++++++++++++++++
.../BasicDataProfilingDataSamplerTest.java | 108 +++++++++++++++
.../neo4j/execution/NeoExecutionInfoLocation.java | 17 ++-
7 files changed, 459 insertions(+), 39 deletions(-)
diff --git
a/core/src/main/java/org/apache/hop/core/row/value/ValueMetaAvroRecord.java
b/core/src/main/java/org/apache/hop/core/row/value/ValueMetaAvroRecord.java
index 39d66ef5d5..6ee5deac50 100644
--- a/core/src/main/java/org/apache/hop/core/row/value/ValueMetaAvroRecord.java
+++ b/core/src/main/java/org/apache/hop/core/row/value/ValueMetaAvroRecord.java
@@ -17,6 +17,8 @@
package org.apache.hop.core.row.value;
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
import java.io.DataInputStream;
import java.io.DataOutputStream;
import java.io.EOFException;
@@ -57,6 +59,11 @@ public class ValueMetaAvroRecord extends ValueMetaBase {
private static final String CONST_SCHEMA = "schema";
private static final String CONST_SPECIFIED = " specified.";
private static final String CONST_UNKNOWN_STORAGE_TYPE = " : Unknown storage
type ";
+ private static final String CONST_SCHEMA_NEEDED =
+ "An Avro schema is needed to read a GenericRecord from an input stream";
+
+ /** Reject a corrupt length instead of allocating it. Sample rows are far
smaller than this. */
+ private static final int MAX_AVRO_BYTES = 32 * 1024 * 1024;
public ValueMetaAvroRecord() {
super(null, IValueMeta.TYPE_AVRO);
@@ -417,12 +424,9 @@ public class ValueMetaAvroRecord extends ValueMetaBase {
outputStream.writeBoolean(object == null);
if (object != null) {
- GenericRecord genericRecord = (GenericRecord) object;
-
- BinaryEncoder binaryEncoder =
EncoderFactory.get().directBinaryEncoder(outputStream, null);
- GenericDatumWriter<GenericRecord> datumWriter =
- new GenericDatumWriter<>(genericRecord.getSchema());
- datumWriter.write(genericRecord, binaryEncoder);
+ // Schema and datum are length-prefixed so the rest of the row stays
aligned, and so a
+ // field with no schema on its metadata (Avro File Input, Kafka
Consumer) can be read back.
+ outputStream.write(encodeRecord(object));
}
} catch (IOException e) {
throw new HopFileException(this + " : Unable to write value data to
output stream", e);
@@ -438,20 +442,14 @@ public class ValueMetaAvroRecord extends ValueMetaBase {
return null; // done
}
- // De-serialize a GenericRow object
- //
- if (schema == null) {
- throw new HopFileException(
- "An Avro schema is needed to read a GenericRecord from an input
stream");
- }
-
- BinaryDecoder binaryDecoder =
DecoderFactory.get().directBinaryDecoder(inputStream, null);
- GenericDatumReader<GenericRecord> datumReader = new
GenericDatumReader<>(schema);
-
- return datumReader.read(null, binaryDecoder);
+ return decodeRecord(inputStream);
+ } catch (HopEofException e) {
+ throw e;
+ } catch (SocketTimeoutException e) {
+ throw e;
} catch (EOFException e) {
throw new HopEofException(e);
- } catch (SocketTimeoutException e) {
+ } catch (HopFileException e) {
throw e;
} catch (IOException e) {
throw new HopFileException(this + " : Unable to read value data from
input stream", e);
@@ -459,8 +457,8 @@ public class ValueMetaAvroRecord extends ValueMetaBase {
}
/**
- * Bytes of the compact schema JSON plus the binary datum. Sample storage
uses this as the size of
- * an Avro value: the schema is part of the payload and {@link #writeData}
writes the datum.
+ * Bytes of the compact schema JSON plus the binary datum. An execution data
profile uses this as
+ * the size of an Avro sample value.
*
* @param object An Avro {@link GenericRecord}
* @return The schema and datum size in bytes
@@ -514,6 +512,117 @@ public class ValueMetaAvroRecord extends ValueMetaBase {
return super.getInteger(object);
}
+ /**
+ * Minimum and maximum profiling, and row-buffer equality, call this. Avro
has no single ordering,
+ * so the record text already shown in the execution grid is used.
+ */
+ @Override
+ protected int typeCompare(Object data1, Object data2) throws
HopValueException {
+ String one = getString(data1);
+ String two = getString(data2);
+ if (one == null && two == null) {
+ return 0;
+ }
+ if (one == null) {
+ return -1;
+ }
+ if (two == null) {
+ return 1;
+ }
+ return one.compareTo(two);
+ }
+
+ @Override
+ public boolean requiresRealClone() {
+ return true;
+ }
+
+ /**
+ * Encode a record so it can be stored without a schema on the value
metadata. The bytes are the
+ * schema JSON (4-byte length, UTF-8) followed by the Avro binary datum
(4-byte length).
+ */
+ public static byte[] encodeRecord(Object object) throws HopFileException {
+ GenericRecord genericRecord = toGenericRecord(object);
+ Schema recordSchema = genericRecord.getSchema();
+ if (recordSchema == null) {
+ throw new HopFileException(CONST_SCHEMA_NEEDED);
+ }
+ try {
+ ByteArrayOutputStream baos = new ByteArrayOutputStream();
+ DataOutputStream dos = new DataOutputStream(baos);
+ byte[] schemaBytes =
recordSchema.toString(false).getBytes(StandardCharsets.UTF_8);
+ dos.writeInt(schemaBytes.length);
+ dos.write(schemaBytes);
+
+ ByteArrayOutputStream avroBytes = new ByteArrayOutputStream();
+ BinaryEncoder binaryEncoder =
EncoderFactory.get().directBinaryEncoder(avroBytes, null);
+ new GenericDatumWriter<GenericRecord>(recordSchema).write(genericRecord,
binaryEncoder);
+ binaryEncoder.flush();
+ byte[] data = avroBytes.toByteArray();
+ dos.writeInt(data.length);
+ dos.write(data);
+ dos.flush();
+ return baos.toByteArray();
+ } catch (Exception e) {
+ throw new HopFileException("Unable to encode an Avro record", e);
+ }
+ }
+
+ /** Decode a payload produced by {@link #encodeRecord(Object)}. */
+ public static GenericRecord decodeRecord(byte[] payload) throws
HopFileException {
+ if (payload == null) {
+ return null;
+ }
+ try (DataInputStream dis = new DataInputStream(new
ByteArrayInputStream(payload))) {
+ return decodeRecord(dis);
+ } catch (HopFileException e) {
+ throw e;
+ } catch (IOException e) {
+ throw new HopFileException("Unable to decode an Avro record", e);
+ }
+ }
+
+ private static GenericRecord toGenericRecord(Object object) throws
HopFileException {
+ if (object instanceof GenericRecord genericRecord) {
+ return genericRecord;
+ }
+ throw new HopFileException(
+ "Expected an Avro GenericRecord and got "
+ + (object == null ? "null" : object.getClass().getName()));
+ }
+
+ private static GenericRecord decodeRecord(DataInputStream inputStream)
throws HopFileException {
+ try {
+ byte[] schemaBytes = readBounded(inputStream);
+ String schemaJson = new String(schemaBytes, StandardCharsets.UTF_8);
+ if (StringUtils.isEmpty(schemaJson)) {
+ throw new HopFileException(CONST_SCHEMA_NEEDED);
+ }
+ Schema recordSchema = new Schema.Parser().parse(schemaJson);
+ byte[] data = readBounded(inputStream);
+ BinaryDecoder binaryDecoder =
+ DecoderFactory.get().binaryDecoder(new ByteArrayInputStream(data),
null);
+ return new GenericDatumReader<GenericRecord>(recordSchema).read(null,
binaryDecoder);
+ } catch (HopFileException e) {
+ throw e;
+ } catch (EOFException e) {
+ throw new HopEofException(e);
+ } catch (Exception e) {
+ throw new HopFileException("Unable to decode an Avro record", e);
+ }
+ }
+
+ private static byte[] readBounded(DataInputStream inputStream)
+ throws IOException, HopFileException {
+ int length = inputStream.readInt();
+ if (length < 0 || length > MAX_AVRO_BYTES) {
+ throw new HopFileException("Avro value length " + length + " is not
valid");
+ }
+ byte[] bytes = new byte[length];
+ inputStream.readFully(bytes);
+ return bytes;
+ }
+
/** Counts bytes written by an Avro encoder without keeping them. */
private static final class CountingOutputStream extends OutputStream {
private long count;
diff --git
a/core/src/test/java/org/apache/hop/core/row/value/ValueMetaAvroRecordTest.java
b/core/src/test/java/org/apache/hop/core/row/value/ValueMetaAvroRecordTest.java
index 3db6b47748..56590fa760 100644
---
a/core/src/test/java/org/apache/hop/core/row/value/ValueMetaAvroRecordTest.java
+++
b/core/src/test/java/org/apache/hop/core/row/value/ValueMetaAvroRecordTest.java
@@ -20,6 +20,7 @@ package org.apache.hop.core.row.value;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import com.fasterxml.jackson.databind.JsonNode;
@@ -33,6 +34,8 @@ import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.util.Utf8;
import org.apache.hop.core.HopClientEnvironment;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
import org.apache.hop.core.util.TestUtil;
import org.json.simple.JSONObject;
import org.json.simple.parser.JSONParser;
@@ -215,6 +218,50 @@ class ValueMetaAvroRecordTest {
verifyGenericRecords(genericRecord, verify);
}
+ @Test
+ void testWriteReadDataWithoutSchemaOnMetadata() throws Exception {
+ GenericRecord genericRecord = generateGenericRecord();
+ ValueMetaAvroRecord valueMeta = new ValueMetaAvroRecord("test");
+
+ ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
+ DataOutputStream outputStream = new
DataOutputStream(byteArrayOutputStream);
+ valueMeta.writeData(outputStream, genericRecord);
+ valueMeta.writeData(outputStream, null);
+ outputStream.close();
+
+ ByteArrayInputStream byteArrayInputStream =
+ new ByteArrayInputStream(byteArrayOutputStream.toByteArray());
+ DataInputStream inputStream = new DataInputStream(byteArrayInputStream);
+ GenericRecord verify = (GenericRecord) valueMeta.readData(inputStream);
+ assertNull(valueMeta.readData(inputStream));
+
+ verifyGenericRecords(genericRecord, verify);
+ assertNull(valueMeta.getSchema());
+ }
+
+ @Test
+ void testWriteReadRowAroundAvroField() throws Exception {
+ GenericRecord genericRecord = generateGenericRecord();
+ IRowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("name"));
+ rowMeta.addValueMeta(new ValueMetaAvroRecord("record"));
+ rowMeta.addValueMeta(new ValueMetaInteger("n"));
+ Object[] row = new Object[] {"hop", genericRecord, 5L};
+
+ ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
+ DataOutputStream outputStream = new
DataOutputStream(byteArrayOutputStream);
+ rowMeta.writeData(outputStream, row);
+ outputStream.close();
+
+ DataInputStream inputStream =
+ new DataInputStream(new
ByteArrayInputStream(byteArrayOutputStream.toByteArray()));
+ Object[] loaded = rowMeta.readData(inputStream);
+
+ assertEquals("hop", loaded[0]);
+ verifyGenericRecords(genericRecord, (GenericRecord) loaded[1]);
+ assertEquals(5L, loaded[2]);
+ }
+
private GenericRecord generateGenericRecord() {
Schema schema = new Schema.Parser().parse(SCHEMA_JSON);
GenericRecord genericRecord = new GenericData.Record(schema);
diff --git a/engine/src/main/java/org/apache/hop/execution/ExecutionData.java
b/engine/src/main/java/org/apache/hop/execution/ExecutionData.java
index 199ffe417a..263aab90a6 100644
--- a/engine/src/main/java/org/apache/hop/execution/ExecutionData.java
+++ b/engine/src/main/java/org/apache/hop/execution/ExecutionData.java
@@ -40,6 +40,7 @@ import lombok.Getter;
import lombok.Setter;
import org.apache.hop.core.Const;
import org.apache.hop.core.exception.HopFileException;
+import org.apache.hop.core.logging.LogChannel;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.row.RowBuffer;
import org.apache.hop.core.row.RowMeta;
@@ -256,26 +257,49 @@ public class ExecutionData {
int nrSets = dis.readInt();
for (int i = 0; i < nrSets; i++) {
- // The set key & description
- //
- String setKey = dis.readUTF();
-
- // The row metadata...
- //
- IRowMeta rowMeta = new RowMeta(dis);
-
- // How many data rows does this buffer have?
- //
- List<Object[]> rows = new ArrayList<>();
- int nrRows = dis.readInt();
- for (int r = 0; r < nrRows; r++) {
- Object[] row = rowMeta.readData(dis);
- rows.add(row);
+ if (!readDataSet(dis)) {
+ // An undecodable value (Avro stored without its schema) must not
reject the cache
+ // entry that contains this blob. Complete rows stay; the rest of
the blob is dropped.
+ return;
}
-
- dataSets.put(setKey, new RowBuffer(rowMeta, rows));
}
}
}
}
+
+ /**
+ * @return false when a value could not be read and the remainder of the
blob was left unread
+ */
+ private boolean readDataSet(DataInputStream dis) throws IOException {
+ String setKey;
+ IRowMeta rowMeta;
+ int nrRows;
+ try {
+ setKey = dis.readUTF();
+ rowMeta = new RowMeta(dis);
+ nrRows = dis.readInt();
+ } catch (HopFileException e) {
+ logStoppedReading(null, e);
+ return false;
+ }
+
+ List<Object[]> rows = new ArrayList<>();
+ for (int r = 0; r < nrRows; r++) {
+ try {
+ rows.add(rowMeta.readData(dis));
+ } catch (HopFileException e) {
+ logStoppedReading(setKey, e);
+ dataSets.put(setKey, new RowBuffer(rowMeta, rows));
+ return false;
+ }
+ }
+ dataSets.put(setKey, new RowBuffer(rowMeta, rows));
+ return true;
+ }
+
+ private static void logStoppedReading(String setKey, HopFileException e) {
+ String where = setKey == null ? "" : " at set '" + setKey + "'";
+ LogChannel.GENERAL.logError(
+ "Stopped reading execution data" + where + ". Rows after this point
were dropped.", e);
+ }
}
diff --git
a/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSampler.java
b/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSampler.java
index 7ad676ab9a..6848e152a2 100644
---
a/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSampler.java
+++
b/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSampler.java
@@ -373,7 +373,7 @@ public class BasicDataProfilingDataSampler
RowBuffer rowBuffer =
typeBufferMap.computeIfAbsent(profilingType, k -> new
RowBuffer(rowMeta));
- // Keep the memory consumption sane
+ // Keep the memory consumption sane. Copy the row: the transform may
reuse the array.
//
if (rowBuffer.size() < store.getMaxRows()) {
rowBuffer.addRow(getSampledValueLimits().copyRow(rowMeta, row,
decisions));
diff --git
a/engine/src/test/java/org/apache/hop/execution/ExecutionDataTest.java
b/engine/src/test/java/org/apache/hop/execution/ExecutionDataTest.java
index a87ca1d62c..c4090567fd 100644
--- a/engine/src/test/java/org/apache/hop/execution/ExecutionDataTest.java
+++ b/engine/src/test/java/org/apache/hop/execution/ExecutionDataTest.java
@@ -20,19 +20,36 @@ package org.apache.hop.execution;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.node.ObjectNode;
+import java.io.ByteArrayOutputStream;
+import java.io.DataOutputStream;
import java.math.BigDecimal;
import java.util.Arrays;
+import java.util.Base64;
import java.util.Date;
import java.util.List;
import java.util.Map;
+import java.util.zip.GZIPOutputStream;
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericDatumWriter;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.io.EncoderFactory;
+import org.apache.avro.util.Utf8;
import org.apache.hop.core.HopClientEnvironment;
import org.apache.hop.core.json.HopJson;
import org.apache.hop.core.row.IRowMeta;
import org.apache.hop.core.row.RowBuffer;
+import org.apache.hop.core.row.RowMeta;
import org.apache.hop.core.row.RowMetaBuilder;
+import org.apache.hop.core.row.value.ValueMetaAvroRecord;
+import org.apache.hop.core.row.value.ValueMetaInteger;
+import org.apache.hop.core.row.value.ValueMetaString;
import org.apache.hop.core.util.TestUtil;
+import org.apache.hop.execution.caching.CacheEntry;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -130,4 +147,104 @@ class ExecutionDataTest {
assertEquals(data, copy);
}
+
+ @Test
+ void testAvroRecordRoundTripWithoutSchemaOnMetadata() throws Exception {
+ Schema schema =
+ new Schema.Parser()
+ .parse(
+
"{\"type\":\"record\",\"name\":\"msg\",\"fields\":[{\"name\":\"body\",\"type\":\"string\"}]}");
+ GenericRecord record = new GenericData.Record(schema);
+ record.put("body", new Utf8("hello"));
+
+ IRowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("name"));
+ rowMeta.addValueMeta(new ValueMetaAvroRecord("record"));
+ rowMeta.addValueMeta(new ValueMetaInteger("n"));
+ List<Object[]> rows = List.<Object[]>of(new Object[] {"hop", record, 5L});
+
+ ExecutionData data =
+ ExecutionDataBuilder.of()
+ .addDataSets(Map.of("rows", new RowBuffer(rowMeta, rows)))
+ .withParentId("parentId")
+ .withOwnerId(ExecutionDataBuilder.ALL_TRANSFORMS)
+ .build();
+
+ ObjectMapper objectMapper = HopJson.newMapper();
+ ExecutionData copy =
+ objectMapper.readValue(objectMapper.writeValueAsString(data),
ExecutionData.class);
+
+ assertEquals(data, copy);
+ GenericRecord loaded = (GenericRecord)
copy.getDataSets().get("rows").getBuffer().get(0)[1];
+ assertEquals(new Utf8("hello"), loaded.get("body"));
+ }
+
+ @Test
+ void testOldAvroPayloadDoesNotRejectCacheEntry() throws Exception {
+ IRowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("name"));
+ rowMeta.addValueMeta(new ValueMetaAvroRecord("record"));
+
+ Schema schema =
+ new Schema.Parser()
+ .parse(
+
"{\"type\":\"record\",\"name\":\"msg\",\"fields\":[{\"name\":\"body\",\"type\":\"string\"}]}");
+ GenericRecord record = new GenericData.Record(schema);
+ record.put("body", "stored-without-schema");
+
+ ByteArrayOutputStream baos = new ByteArrayOutputStream();
+ try (GZIPOutputStream gzos = new GZIPOutputStream(baos);
+ DataOutputStream dos = new DataOutputStream(gzos)) {
+ dos.writeInt(1);
+ dos.writeUTF("samples");
+ rowMeta.writeMeta(dos);
+ dos.writeInt(2);
+ // A null Avro value has the same layout as before: a single true flag.
+ rowMeta.writeData(dos, new Object[] {"before", null});
+ // The old writer then streamed a raw Avro encoder with no length and no
schema.
+ new ValueMetaString("name").writeData(dos, "after");
+ dos.writeBoolean(false);
+ GenericDatumWriter<GenericRecord> writer = new
GenericDatumWriter<>(record.getSchema());
+ writer.write(record, EncoderFactory.get().directBinaryEncoder(dos,
null));
+ }
+ String oldBlob = Base64.getEncoder().encodeToString(baos.toByteArray());
+
+ IRowMeta keptMeta = new RowMetaBuilder().addString("name").build();
+ ExecutionData sibling =
+ ExecutionDataBuilder.of()
+ .withOwnerId("transform-copy")
+ .withParentId("parent-id")
+ .addDataSets(
+ Map.of("rows", new RowBuffer(keptMeta, List.<Object[]>of(new
Object[] {"kept"}))))
+ .build();
+
+ CacheEntry entry = new CacheEntry();
+ entry.setId("parent-id");
+ entry.setName("kafka-pipeline");
+ entry.getChildExecutionData().put("transform-copy", sibling);
+
+ ObjectMapper objectMapper = HopJson.newMapper();
+ ObjectNode node = (ObjectNode)
objectMapper.readTree(objectMapper.writeValueAsString(entry));
+ ObjectNode allTransforms = objectMapper.createObjectNode();
+ allTransforms.put("ownerId", ExecutionDataBuilder.ALL_TRANSFORMS);
+ allTransforms.put("parentId", "parent-id");
+ allTransforms.put("rowsBinaryGzipBase64Encoded", oldBlob);
+ ((ObjectNode) node.get("childExecutionData"))
+ .set(ExecutionDataBuilder.ALL_TRANSFORMS, allTransforms);
+
+ CacheEntry loaded =
+ objectMapper.readValue(objectMapper.writeValueAsString(node),
CacheEntry.class);
+
+ assertEquals("parent-id", loaded.getId());
+ assertEquals("kafka-pipeline", loaded.getName());
+ ExecutionData loadedSibling =
loaded.getChildExecutionData().get("transform-copy");
+ assertEquals("kept",
loadedSibling.getDataSets().get("rows").getBuffer().get(0)[0]);
+
+ ExecutionData samples =
loaded.getChildExecutionData().get(ExecutionDataBuilder.ALL_TRANSFORMS);
+ assertNotNull(samples);
+ List<Object[]> keptRows = samples.getDataSets().get("samples").getBuffer();
+ assertEquals(1, keptRows.size());
+ assertEquals("before", keptRows.get(0)[0]);
+ assertNull(keptRows.get(0)[1]);
+ }
}
diff --git
a/engine/src/test/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerTest.java
b/engine/src/test/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerTest.java
new file mode 100644
index 0000000000..4298a8aa75
--- /dev/null
+++
b/engine/src/test/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerTest.java
@@ -0,0 +1,108 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.execution.sampler.plugins.dataprof;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.util.Utf8;
+import org.apache.hop.core.HopClientEnvironment;
+import org.apache.hop.core.json.HopJson;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowBuffer;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaAvroRecord;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.core.util.TestUtil;
+import org.apache.hop.core.variables.Variables;
+import org.apache.hop.execution.ExecutionData;
+import org.apache.hop.execution.ExecutionDataBuilder;
+import org.apache.hop.execution.sampler.ExecutionDataSamplerMeta;
+import
org.apache.hop.execution.sampler.plugins.dataprof.BasicDataProfilingDataSampler.ProfilingType;
+import org.apache.hop.pipeline.transform.stream.IStream;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+class BasicDataProfilingDataSamplerTest {
+
+ @BeforeAll
+ static void before() throws Exception {
+ HopClientEnvironment.init();
+ TestUtil.registerTestPluginTypes();
+ }
+
+ @Test
+ void avroColumnIsProfiledAndSurvivesExecutionData() throws Exception {
+ BasicDataProfilingDataSampler sampler = new
BasicDataProfilingDataSampler();
+ sampler.setSampleSize("10");
+ sampler.setOnlyProfilingLastTransforms(true);
+ ExecutionDataSamplerMeta samplerMeta =
+ new ExecutionDataSamplerMeta("Kafka consumer", "0", "log-channel",
false, true);
+ BasicDataProfilingDataSamplerStore store =
sampler.createSamplerStore(samplerMeta);
+ store.init(new Variables(), null, null);
+
+ IRowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("id"));
+ rowMeta.addValueMeta(new ValueMetaAvroRecord("message"));
+
+ GenericRecord first = record("alpha");
+ GenericRecord second = record("beta");
+ sampler.sampleRow(store, IStream.StreamType.OUTPUT, rowMeta, new Object[]
{"1", first});
+ sampler.sampleRow(store, IStream.StreamType.OUTPUT, rowMeta, new Object[]
{"2", second});
+
+ assertNotNull(store.getMinValues().get("message"));
+ assertNotNull(store.getMaxValues().get("message"));
+
+ sampler.sampleRow(store, IStream.StreamType.OUTPUT, rowMeta, new Object[]
{"3", null});
+
+ assertEquals(1L, store.getNullCounters().get("message"));
+ assertEquals(2L, store.getNonNullCounters().get("message"));
+ assertNotNull(store.getMaxValues().get("message"));
+
+ RowBuffer nonNullSamples =
+ store.getProfileSamples().get("message").get(ProfilingType.NrNonNulls);
+ assertNotSame(first, nonNullSamples.getBuffer().get(0)[1]);
+
+ ExecutionData data =
+ ExecutionDataBuilder.of()
+ .withParentId("parent")
+ .withOwnerId(ExecutionDataBuilder.ALL_TRANSFORMS)
+ .addDataSets(store.getSamples())
+ .addSetMeta(store.getSamplesMetadata())
+ .build();
+
+ ObjectMapper mapper = HopJson.newMapper();
+ ExecutionData copy = mapper.readValue(mapper.writeValueAsString(data),
ExecutionData.class);
+ assertEquals(data, copy);
+ }
+
+ private static GenericRecord record(String body) {
+ Schema schema =
+ new Schema.Parser()
+ .parse(
+
"{\"type\":\"record\",\"name\":\"msg\",\"fields\":[{\"name\":\"body\",\"type\":\"string\"}]}");
+ GenericRecord record = new GenericData.Record(schema);
+ record.put("body", new Utf8(body));
+ return record;
+ }
+}
diff --git
a/plugins/tech/neo4j/src/main/java/org/apache/hop/neo4j/execution/NeoExecutionInfoLocation.java
b/plugins/tech/neo4j/src/main/java/org/apache/hop/neo4j/execution/NeoExecutionInfoLocation.java
index ddc6516d2a..d09bce453b 100644
---
a/plugins/tech/neo4j/src/main/java/org/apache/hop/neo4j/execution/NeoExecutionInfoLocation.java
+++
b/plugins/tech/neo4j/src/main/java/org/apache/hop/neo4j/execution/NeoExecutionInfoLocation.java
@@ -48,6 +48,7 @@ import org.apache.hop.core.row.IValueMeta;
import org.apache.hop.core.row.JsonRowMeta;
import org.apache.hop.core.row.RowBuffer;
import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaAvroRecord;
import org.apache.hop.core.row.value.ValueMetaJson;
import org.apache.hop.core.variables.IVariables;
import org.apache.hop.execution.Execution;
@@ -1373,7 +1374,13 @@ public class NeoExecutionInfoLocation implements
IExecutionInfoLocation {
IValueMeta valueMeta = rowMeta.getValueMeta(v);
Object valueData = null;
try {
- valueData = valueMeta.getNativeDataType(row[v]);
+ if (valueMeta.getType() == IValueMeta.TYPE_AVRO) {
+ // A GenericRecord is not a Neo4j property. Store the same
length-prefixed bytes the
+ // row codec writes into an execution-data blob.
+ valueData = row[v] == null ? null :
ValueMetaAvroRecord.encodeRecord(row[v]);
+ } else {
+ valueData = valueMeta.getNativeDataType(row[v]);
+ }
} catch (Exception e) {
if (nrErrors++ < 10) {
log.logError(
@@ -1641,6 +1648,14 @@ public class NeoExecutionInfoLocation implements
IExecutionInfoLocation {
yield e.getMessage();
}
}
+ case IValueMeta.TYPE_AVRO -> {
+ try {
+ yield ValueMetaAvroRecord.decodeRecord(value.asByteArray());
+ } catch (Exception e) {
+ throw new HopRuntimeException(
+ "Unable to read Avro value '" + valueMeta.getName() + "'", e);
+ }
+ }
default ->
// Convert from String
//