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
       //

Reply via email to