Hi Jark,

Thanks for the thorough review! I've addressed all 9 points in the updated 
design document:

1. Read Behavior on Type Mismatch: Added explicit definition — when typed_value 
is null but value is non-null during column pruning, the reader returns null. 
Clients that need the actual value must also project the value (binary 
fallback) column.

2. Unified Project API: Removed variantFieldProjection(Map<Integer, 
List<String>>). Variant field projection is now unified into the .project() API 
with syntax "variant_col:field.path::type" (e.g., "payload:action::STRING", 
"payload:user.name::STRING").

3. Tiering Integration: Expanded the Stream-Lake Integration section to 
explicitly cover Paimon (supported in this FIP), Iceberg (future work), and 
Lance (future work) with scope clarification.

4. Adapter Pattern instead of Reflection: Replaced all reflection-based type 
mapping with an Adapter pattern approach, referencing SchemaAdapter as 
reference implementation for Flink 2.x.

5. @PublicEvolving: Changed all new API annotations from @PublicStable to 
@PublicEvolving.

6. Table Property Prefix + Example: Renamed to 
table.variant.shredding.schema.<column_name> and added a concrete JSON 
configuration example.

7. Client Writer Prefix: All writer options now use 
client.writer.variant.shredding.* prefix.

8. Arrow Java 19.0: Updated to recommend considering the upgrade, noting Java 8 
compatibility is no longer a constraint since Fluss only publishes Java 11 
artifacts. I will open a separate PR to drive the Arrow Java upgrade.

9. Spark Integration: Added a new "6. Spark Connector Integration" section 
covering type mapping, row data conversion, variant field pushdown, and version 
compatibility.

Please take another look when you get a chance.

Best,
Yunhong (Swuferhong)

On 2026/05/16 13:21:47 Jark Wu wrote:
> Thanks for the detailed design document.
> 
> I have several questions and suggestions:
> 
> 1. Question: Behavior When Reading Mismatched Types from Shredded Columns
> 
> When both typed_value (e.g., INT) and value (type mismatch fallback)
> exist for a shredded field, and the client performs column pruning to
> read this field as INT, what happens when the reader encounters a row
> where typed_value is null but value is non-null (indicating the actual
> value doesn't match the shredded type)? This needs to be clearly
> defined. Null or exception?
> 
> 2. Client API variantFieldProjection Is Too Complex
> The current API: ".variantFieldProjection(Map<Integer, List<String>>
> columnToFields)"
> is hard to understand from a user's perspective. What is "column
> index"? How do you express nested field access like col.nest1.a? This
> API puts too much burden on the user to understand internal column
> indexing.
> 
> My suggestion: Variant column pruning should be unified into the
> ".project() " API, just like SQL's SELECT syntax handles projection
> uniformly. We can design a special access pattern for variant fields,
> for example:
> "variant_col:a.b.c::int"
> "variant_col:arr[1].field::double"
> 
> 3. Tiering Integration
> The document only elaborates on Paimon integration. What about Iceberg
> and Lance? Can they also receive the tiered variant data? If not
> supported in this FIP, please clarify the scope and mention them as
> future work. If they are planned, please describe the integration
> approach.
> 
> 4. Do Not Use Reflection for Flink Type Mapping
> "Both mappings use reflection to access Flink's VariantType class"
> 
> Please do not use reflection. Use an Adapter pattern instead, refer to
> SchemaAdapter as a reference implementation which API is only
> supported in Flink 2.x.
> 
> 5. @PublicStable Annotation Is Premature
> All new APIs should be marked as @PublicEvolving initially. We can
> upgrade them to @PublicStable after at least 3 minor versions of
> proven stability. This gives us room to iterate on the API without
> breaking compatibility guarantees.
> 
> 6. Clarify "variant.shredding.schema.<column_name>" Parameter
> What is this parameter used for? The document should provide a
> concrete example showing how to configure it.
> Also, if this is a table storage property, it should follow the
> existing convention and use the table.* prefix, e.g.,
> table.variant.shredding.schema.<column_name>.
> 
> 7. Writer Option Prefix Should Be "client.writer.*"
> All the "writer.variant.shredding.*" options should be prefixed with
> "client.writer.*" to align with other writer options in Fluss.
> 
> 8. Consider Upgrading to Arrow Java 19.0 for Variant Support
> The FIP mentions that upgrading to arrow-java 19.0 is not required,
> but I suggest we consider it. Could adopting arrow-java 19.0 simplify
> our implementation, particularly regarding Arrow Variant type
> definitions and read/write paths? Since Fluss currently only publishes
> Java 11 artifacts, we can drop Java 8 compatibility and upgrade to the
> latest version of arrow-java. If this upgrade brings development
> simplifications, I believe it is worth considering.
> 
> 9. Spark Integration Approach
> Spark is another important connector in the code base. I suggest to
> add a Spark section about how to integrate with Spark Variant type,
> and maybe also mention the nested variant type pushdown.
> 
> Best,
> Jark
> 
> On Thu, 7 May 2026 at 09:58, Yunhong Zheng <[email protected]> wrote:
> >
> > Hi Anton,
> >
> > Really thanks for the detailed review. Addressing each point:
> >
> > 1. Writer-side vs. Server-coordinated
> > Writer-independent is our chose. As the FIP states:
> >   - Simplest implementation: no RPC, no persistence, no multi-Writer 
> > coordination
> >   - No single point of failure: Writer is fully autonomous
> >   - Fast startup: new Writer starts shredding after local sampling, no 
> > server dependency
> >
> > For multi-language bindings: Fluss's multi-language clients share a unified 
> > Rust core (Python via PyO3, C++ via FFI), so the policy layer (inferrer, 
> > statistics, thresholds) only needs one Rust port — not N independent 
> > implementations. We don't plan to pursue server-coordinated in the near 
> > term, as client-side offers more flexibility.
> >
> > 2. Lance
> > Correct — fluss-lake-lance does not yet handle Variant shredding. Lance 
> > upstream Variant support is still in progress. We'll document this as an 
> > explicit limitation and gate it behind a compatibility check.
> >
> > 3. Arrow Java Migration
> > Upgrading to Arrow Java 19.0.0 requires significant changes on the Fluss 
> > side. Since we only use StructVector for the shredding layout, it's 
> > practical to implement independently in Fluss for now. Once Arrow Java 
> > ships more complete shredding support, we can do a smooth integration — the 
> > current implementation is designed to be convergence-friendly.
> >
> > 4. Doc/Code Alignment
> > This has been resolved in the current FIP — Section 4.4 explicitly commits 
> > to the client-side independent computation approach.
> >
> > 5. WAL buildAndWriteResidual per-row cost
> > Acknowledged — v1 prioritizes correctness over allocation efficiency. 
> > Post-v1 targets: buffer pooling, batch-level residual encoding, and 
> > incremental encoding for UPDATE. Will add a note in "Performance 
> > Considerations" to make this explicit.
> >
> > Best,
> > Yunhong Zheng (Swuferhong)
> >
> > On 2026/05/02 02:33:04 Anton Borisov wrote:
> > > Hi Yunhong,
> > >
> > > +1 directionally. I read through both the FIP and #2981, and had a few
> > > questions and comments.
> > >
> > > Bindings: arrow-rs's parquet-variant covers much of the spec
> > > mechanics, so the encode/decode/shred layer has a reasonable
> > > upstream path. What each binding would still need to port is the
> > > policy layer Fluss adds on top: inferrer, statistics collector,
> > > thresholds, candidate ranking, and tie-breaking. With shredding
> > > decisions made per writer (as the writer.variant.shredding.*
> > > configuration implies), each port has to produce compatible
> > > decisions with the Java inferrer for the same input, otherwise
> > > batches may fragment further. The doc in the PR itself acknowledges
> > > this can produce heterogeneous batches across Java writers,
> > > multi-language writers amplify that. A server-coordinated alternative
> > >  would sidestep most of this. Is writer-side the v1 commitment, or
> > > Is server-coordination on the post-1.0 roadmap?
> > > Asking because the answer changes the binding strategy.
> > >
> > > Lance: the PR covers Iceberg and Paimon, but I did not see
> > > corresponding changes for fluss-lake-lance. Lance upstream Variant
> > > support also still appears to be in progress, and
> > > ArrowDataConverter.convertToNonShaded does not seem to handle
> > > StructVector. If a user tiers Variant data to Lance today, what
> > > happens? We might want to be explicit here.
> > >
> > > Arrow Java migration: the Java side implements custom shredding
> > > because Arrow Java does not yet provide native shredding support.
> > > When upstream catches up, do we intend to converge on Arrow Java's
> > > implementation, or keep the Fluss implementation as the canonical
> > > one for stability?
> > >
> > > Doc/code alignment: the design doc weighs writer-independent vs
> > > server-coordinated shredding as still open, but
> > > VariantShreddingManager and the writer.variant.shredding.* config
> > > commit to writer-independent.
> > >
> > > One observation while tracing the WAL path: walBuilder.append
> > > re-encodes Variant via ArrowShreddedVariantWriter, so UPDATE_BEFORE
> > > on an old-layout row gets re-shredded under the current policy.
> > > That keeps BEFORE/AFTER consistent, but buildAndWriteResidual
> > > appears to allocate per row: two ArrayLists, Arrays.copyOfRange per
> > > non-shredded field, and a fresh byte[] from encodeObject. That
> > > means UPDATE pays this twice, while INSERT pays it once. Not a
> > > blocker, but it is a non-obvious write-path cost from reading the FIP.
> > >
> > >  -- Anton
> > >
> > > вт, 28 апр. 2026 г. в 02:29, yunhong Zheng <[email protected]>:
> > > >
> > > > Hi all,
> > > >
> > > > Semi-structured data (JSON, event payloads, metadata bags) is
> > > > ubiquitous in modern data pipelines. In the current Fluss
> > > > implementation, users are forced to store such data in STRING columns
> > > > and parse it at query time. This approach suffers from several
> > > > fundamental problems: no compact binary encoding, no field-level
> > > > access, and no ability to leverage Fluss's existing columnar
> > > > optimizations (column pruning, predicate pushdown).
> > > >
> > > > Meanwhile, the industry is converging on the Parquet Variant Binary
> > > > Encoding specification as the standard for semi-structured data —
> > > > Apache Spark 4.0, Flink 2.1 (FLIP-521), Paimon, and Parquet itself are
> > > > all adopting it. As a streaming storage layer bridging Flink and
> > > > Paimon in stream-lake architectures, Fluss needs native Variant
> > > > support to enable seamless data flow across the ecosystem.
> > > >
> > > > So, I'd like to propose FIP-36: Support Variant Type and Shredding[1].
> > > > This proposal introduces a three-layer design:
> > > >
> > > > 1. Variant Binary Encoding: Adopt the Parquet Variant spec as the
> > > > canonical binary representation, enabling O(log n) field lookup and
> > > > zero-copy nested extraction.
> > > >
> > > > 2. Arrow Columnar Storage: Store Variant values in Arrow StructVector
> > > > with decomposed metadata/value children, integrating with Fluss's
> > > > existing zero-copy RecordBatch infrastructure.
> > > >
> > > > 3. Shredding for Query Optimization: Automatically extract frequently
> > > > accessed fields into typed sub-columns within the Variant
> > > > StructVector, enabling server-side sub-column pruning on
> > > > semi-structured data — all transparent to users.
> > > >
> > > > Any feedback and suggestions on this proposal are welcome!
> > > >
> > > > [1]: 
> > > > https://cwiki.apache.org/confluence/display/FLUSS/FIP-36%3A+Support+Variant+Type+and+Shredding
> > > >
> > > > Regards,
> > > > Yunhong (Swuferhong)
> > >
> 

Reply via email to