This is an automated email from the ASF dual-hosted git repository.
CurtHagenlocher pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-dotnet.git
The following commit(s) were added to refs/heads/main by this push:
new c8aa92f feat: Add array-level shredding entry points on VariantArray
(#447)
c8aa92f is described below
commit c8aa92fb6aa69bbe45b9be162b62eb82f0fd846b
Author: Curt Hagenlocher <[email protected]>
AuthorDate: Sun Sep 27 07:50:35 2026 -0700
feat: Add array-level shredding entry points on VariantArray (#447)
## What's Changed
The shredding APIs only worked one value at a time, so every caller
holding a `VariantArray` had to write the same loops: decode each row
(plus a null mask) to shred it, and rebuild row by row to unshred it.
This adds extension methods on `VariantArray` in
`VariantArrayShreddingExtensions`:
- `Reassemble(allocator = null)` converts a shredded array into its
unshredded equivalent. Null elements stay null, and an unshredded input
is returned unchanged.
- `Shred(schema, allocator = null)` shreds into the given layout. It
reads logical values, so an already-shredded input is reshredded rather
than losing its typed columns.
- `TryShred(options, out shredded, allocator = null)` infers a schema
and shreds into it. It returns false (with `shredded` set to null) when
the inferred schema is unshredded.
- `InferShredSchema(options = null)` infers a schema from an array. This
isn't in the issue's list, but it's needed for the batch scenario the
issue describes: a Parquet writer infers once over a representative
batch, then calls `Shred(schema)` on every batch so all row groups share
one layout.
All four go through one shared loop that resolves the column's schema
and child arrays once rather than per row. Null elements pass through
the `VariantValue?` pipeline added in #445.
Closes #399.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../Shredding/VariantArrayShreddingExtensions.cs | 105 +++++++++++
.../Shredding/VariantArrayShreddingTests.cs | 209 +++++++++++++++++++++
2 files changed, 314 insertions(+)
diff --git
a/src/Apache.Arrow.Operations/Shredding/VariantArrayShreddingExtensions.cs
b/src/Apache.Arrow.Operations/Shredding/VariantArrayShreddingExtensions.cs
index 4e1fbfe..4a685d2 100644
--- a/src/Apache.Arrow.Operations/Shredding/VariantArrayShreddingExtensions.cs
+++ b/src/Apache.Arrow.Operations/Shredding/VariantArrayShreddingExtensions.cs
@@ -14,7 +14,9 @@
// limitations under the License.
using System;
+using System.Collections.Generic;
using Apache.Arrow;
+using Apache.Arrow.Memory;
using Apache.Arrow.Scalars.Variant;
namespace Apache.Arrow.Operations.Shredding
@@ -78,6 +80,109 @@ namespace Apache.Arrow.Operations.Shredding
return GetShreddedVariant(array, index).ToVariantValue();
}
+ /// <summary>
+ /// Infers a <see cref="ShredSchema"/> from the logical values of a
variant array.
+ /// Null elements are ignored. Works for both shredded and unshredded
columns.
+ /// </summary>
+ /// <remarks>
+ /// When writing a column in batches (e.g. Parquet row groups), infer
once over a
+ /// representative batch and pass the result to <see cref="Shred"/>
for every batch,
+ /// so that all batches share one layout.
+ /// </remarks>
+ public static ShredSchema InferShredSchema(this VariantArray array,
ShredOptions options = null)
+ {
+ if (array == null) throw new ArgumentNullException(nameof(array));
+ return new ShredSchemaInferer().Infer(GetLogicalValues(array),
options);
+ }
+
+ /// <summary>
+ /// Shreds a variant array into the layout described by <paramref
name="schema"/>.
+ /// The input may itself be shredded (under any schema); its logical
values are
+ /// re-shredded. Null elements remain null.
+ /// </summary>
+ /// <param name="array">The variant array to shred.</param>
+ /// <param name="schema">The target shredding schema.</param>
+ /// <param name="allocator">Arrow memory allocator, or default if
null.</param>
+ public static VariantArray Shred(this VariantArray array, ShredSchema
schema, MemoryAllocator allocator = null)
+ {
+ if (array == null) throw new ArgumentNullException(nameof(array));
+ if (schema == null) throw new
ArgumentNullException(nameof(schema));
+
+ (byte[] metadata, IReadOnlyList<ShredResult> rows) =
+ VariantShredder.Shred(GetLogicalValues(array), schema);
+ return ShreddedVariantArrayBuilder.Build(schema, metadata, rows,
allocator);
+ }
+
+ /// <summary>
+ /// Infers a shredding schema from <paramref name="array"/> and, if it
produces a
+ /// shredded layout, shreds the array into it.
+ /// </summary>
+ /// <param name="array">The variant array to shred.</param>
+ /// <param name="options">Inference options, or <see
cref="ShredOptions.Default"/> if null.</param>
+ /// <param name="shredded">The shredded array, or null when this
method returns false.</param>
+ /// <param name="allocator">Arrow memory allocator, or default if
null.</param>
+ /// <returns>
+ /// True if a shredded layout was inferred; false if the values have
no layout
+ /// worth shredding (the inferred schema is <see
cref="ShredSchema.Unshredded"/>).
+ /// </returns>
+ public static bool TryShred(
+ this VariantArray array,
+ ShredOptions options,
+ out VariantArray shredded,
+ MemoryAllocator allocator = null)
+ {
+ ShredSchema schema = InferShredSchema(array, options);
+ if (schema.TypedValueType == ShredType.None)
+ {
+ shredded = null;
+ return false;
+ }
+ shredded = Shred(array, schema, allocator);
+ return true;
+ }
+
+ /// <summary>
+ /// Converts a shredded variant array into its unshredded equivalent,
in which
+ /// every element is stored as self-contained metadata and value
bytes. Null
+ /// elements remain null. An unshredded input is returned unchanged.
+ /// </summary>
+ /// <param name="array">The variant array to reassemble.</param>
+ /// <param name="allocator">Arrow memory allocator, or default if
null.</param>
+ public static VariantArray Reassemble(this VariantArray array,
MemoryAllocator allocator = null)
+ {
+ if (array == null) throw new ArgumentNullException(nameof(array));
+ if (!array.IsShredded) return array;
+
+ var builder = new VariantArray.Builder();
+ builder.AppendRange(GetLogicalValues(array));
+ return builder.Build(allocator);
+ }
+
+ /// <summary>
+ /// Enumerates the logical value of every element, with null for null
elements.
+ /// Resolves the column's schema and child arrays once rather than per
row.
+ /// </summary>
+ private static IEnumerable<VariantValue?>
GetLogicalValues(VariantArray array)
+ {
+ ShredSchema schema = GetShredSchema(array);
+ IArrowArray valueArr = array.VariantType.HasValueColumn ?
GetValueArray(array) : null;
+ IArrowArray typedValueArr = array.TypedValueArray;
+
+ for (int i = 0; i < array.Length; i++)
+ {
+ yield return array.IsNull(i)
+ ? (VariantValue?)null
+ : GetLogicalValue(array, schema, valueArr, typedValueArr,
i);
+ }
+ }
+
+ private static VariantValue GetLogicalValue(
+ VariantArray array, ShredSchema schema, IArrowArray valueArr,
IArrowArray typedValueArr, int index)
+ {
+ return new ShreddedVariant(schema, array.GetMetadataBytes(index),
valueArr, typedValueArr, index)
+ .ToVariantValue();
+ }
+
/// <summary>
/// Returns the underlying <c>value</c> sub-array of the
VariantArray's struct storage.
/// This mirrors what <see cref="VariantArray.GetValueBytes"/> uses
internally.
diff --git
a/test/Apache.Arrow.Operations.Tests/Shredding/VariantArrayShreddingTests.cs
b/test/Apache.Arrow.Operations.Tests/Shredding/VariantArrayShreddingTests.cs
new file mode 100644
index 0000000..e637de8
--- /dev/null
+++ b/test/Apache.Arrow.Operations.Tests/Shredding/VariantArrayShreddingTests.cs
@@ -0,0 +1,209 @@
+// 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.
+
+using System.Collections.Generic;
+using Apache.Arrow;
+using Apache.Arrow.Operations.Shredding;
+using Apache.Arrow.Scalars.Variant;
+using Xunit;
+
+namespace Apache.Arrow.Operations.Tests.Shredding
+{
+ /// <summary>
+ /// Tests for the array-level shredding entry points on <see
cref="VariantArray"/>:
+ /// <c>InferShredSchema</c>, <c>Shred</c>, <c>TryShred</c> and
<c>Reassemble</c>.
+ /// </summary>
+ public class VariantArrayShreddingTests
+ {
+ private static VariantValue Obj(int a, string b = null)
+ {
+ var fields = new Dictionary<string, VariantValue> { ["a"] =
VariantValue.FromInt32(a) };
+ if (b != null) fields["b"] = VariantValue.FromString(b);
+ return VariantValue.FromObject(fields);
+ }
+
+ private static VariantArray Unshredded(IEnumerable<VariantValue?>
values)
+ {
+ var builder = new VariantArray.Builder();
+ builder.AppendRange(values);
+ return builder.Build();
+ }
+
+ private static void AssertRows(IReadOnlyList<VariantValue?> expected,
VariantArray array)
+ {
+ Assert.Equal(expected.Count, array.Length);
+ for (int i = 0; i < expected.Count; i++)
+ {
+ if (expected[i].HasValue)
+ {
+ Assert.False(array.IsNull(i), $"row {i} should be valid");
+ Assert.Equal(expected[i].Value,
array.GetLogicalVariantValue(i));
+ }
+ else
+ {
+ Assert.True(array.IsNull(i), $"row {i} should be null");
+ }
+ }
+ }
+
+ private static string Describe(ShredSchema schema)
+ {
+ switch (schema.TypedValueType)
+ {
+ case ShredType.Object:
+ var fields = new List<string>();
+ foreach (KeyValuePair<string, ShredSchema> f in
schema.ObjectFields)
+ fields.Add(f.Key + ":" + Describe(f.Value));
+ fields.Sort(System.StringComparer.Ordinal);
+ return "{" + string.Join(",", fields) + "}";
+ case ShredType.Array:
+ return "[" + Describe(schema.ArrayElement) + "]";
+ default:
+ return schema.TypedValueType.ToString();
+ }
+ }
+
+ private static readonly ShredSchema ObjectA = ShredSchema.ForObject(
+ new Dictionary<string, ShredSchema> { ["a"] =
ShredSchema.Primitive(ShredType.Int32) });
+
+ // Covers a SQL-NULL row, a present variant null, and a value that
falls
+ // back to the residual under an object schema.
+ private static readonly List<VariantValue?> MixedRows = new
List<VariantValue?>
+ {
+ Obj(1, "x"), null, VariantValue.Null, Obj(4),
VariantValue.FromString("not an object"),
+ };
+
+ // Consistent enough to infer an object schema under the default
options.
+ private static readonly List<VariantValue?> ObjectRows = new
List<VariantValue?>
+ {
+ Obj(1, "x"), null, Obj(2), Obj(3, "y"), Obj(4, "z"),
+ };
+
+ [Fact]
+ public void Shred_UnshreddedInput()
+ {
+ VariantArray shredded = Unshredded(MixedRows).Shred(ObjectA);
+
+ Assert.True(shredded.IsShredded);
+ Assert.Equal(Describe(ObjectA),
Describe(shredded.GetShredSchema()));
+ Assert.Equal(1, shredded.NullCount);
+ AssertRows(MixedRows, shredded);
+ }
+
+ [Fact]
+ public void Shred_WithInferredSchema()
+ {
+ VariantArray input = Unshredded(ObjectRows);
+ ShredSchema schema = input.InferShredSchema();
+
+ Assert.Equal("{a:Int32,b:String}", Describe(schema));
+ AssertRows(ObjectRows, input.Shred(schema));
+ }
+
+ [Fact]
+ public void Shred_AlreadyShreddedInput_ReshredsLogicalValues()
+ {
+ VariantArray first = Unshredded(MixedRows).Shred(ObjectA);
+
+ ShredSchema second = ShredSchema.ForObject(
+ new Dictionary<string, ShredSchema> { ["b"] =
ShredSchema.Primitive(ShredType.String) });
+ VariantArray reshredded = first.Shred(second);
+
+ Assert.Equal(Describe(second),
Describe(reshredded.GetShredSchema()));
+ AssertRows(MixedRows, reshredded);
+ }
+
+ [Fact]
+ public void Shred_SlicedInput()
+ {
+ var sliced =
(VariantArray)ArrowArrayFactory.Slice(Unshredded(MixedRows), 1, 3);
+
+ VariantArray shredded = sliced.Shred(ObjectA);
+
+ AssertRows(MixedRows.GetRange(1, 3), shredded);
+ }
+
+ [Fact]
+ public void Shred_BatchesShareInferredLayout()
+ {
+ VariantArray batch1 = Unshredded(new List<VariantValue?> { Obj(1,
"x"), Obj(2, "y") });
+ VariantArray batch2 = Unshredded(new List<VariantValue?> {
VariantValue.FromInt32(7), null });
+ ShredSchema schema = batch1.InferShredSchema();
+
+ VariantArray shredded1 = batch1.Shred(schema);
+ VariantArray shredded2 = batch2.Shred(schema);
+
+ Assert.Equal(Describe(schema),
Describe(shredded1.GetShredSchema()));
+ Assert.Equal(Describe(schema),
Describe(shredded2.GetShredSchema()));
+ AssertRows(new List<VariantValue?> { VariantValue.FromInt32(7),
null }, shredded2);
+ }
+
+ [Fact]
+ public void InferShredSchema_IgnoresNullElements()
+ {
+ VariantArray input = Unshredded(new List<VariantValue?>
+ {
+ VariantValue.FromInt32(1), null, null, null,
VariantValue.FromInt32(2), null,
+ });
+ Assert.Equal(ShredType.Int32,
input.InferShredSchema().TypedValueType);
+ }
+
+ [Fact]
+ public void TryShred_ShreddableValues()
+ {
+ Assert.True(Unshredded(ObjectRows).TryShred(ShredOptions.Default,
out VariantArray shredded));
+ Assert.True(shredded.IsShredded);
+ AssertRows(ObjectRows, shredded);
+ }
+
+ [Fact]
+ public void TryShred_NothingToShred()
+ {
+ VariantArray input = Unshredded(new List<VariantValue?> { null,
null });
+ Assert.False(input.TryShred(null, out VariantArray shredded));
+ Assert.Null(shredded);
+ }
+
+ [Fact]
+ public void Reassemble_ShreddedInput()
+ {
+ VariantArray shredded = Unshredded(MixedRows).Shred(ObjectA);
+
+ VariantArray reassembled = shredded.Reassemble();
+
+ Assert.False(reassembled.IsShredded);
+ Assert.Equal(1, reassembled.NullCount);
+ AssertRows(MixedRows, reassembled);
+ // An unshredded array supports the core (non-Operations) reader.
+ Assert.Equal(MixedRows[0].Value, reassembled.GetVariantValue(0));
+ }
+
+ [Fact]
+ public void Reassemble_SlicedInput()
+ {
+ VariantArray shredded = Unshredded(MixedRows).Shred(ObjectA);
+ var sliced = (VariantArray)ArrowArrayFactory.Slice(shredded, 2, 3);
+
+ AssertRows(MixedRows.GetRange(2, 3), sliced.Reassemble());
+ }
+
+ [Fact]
+ public void Reassemble_UnshreddedInput_ReturnsSameArray()
+ {
+ VariantArray input = Unshredded(MixedRows);
+ Assert.Same(input, input.Reassemble());
+ }
+ }
+}