[ 
https://issues.apache.org/jira/browse/FLINK-40241?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Ramin Gharib updated FLINK-40241:
---------------------------------
    Description: 
*part to* FLINK-37922,  *relates to* FLINK-40218 , FLINK-39605, FLINK-39401

The {{raw}} format is the only way to read a topic or a file that carries no 
schema. VARIANT is the type Flink added for semi-structured data. The two 
cannot be combined today.
{code:sql}
CREATE TABLE events (payload VARIANT) WITH ('format' = 'raw', ...);
{code}
fails at plan time with:
{noformat}
org.apache.flink.table.api.ValidationException:
  The 'raw' format doesn't supports 'VARIANT' as column type.
{noformat}
{{RawFormatFactory#supportedTypes}} lists CHAR, VARCHAR, BINARY, VARBINARY, 
RAW, BOOLEAN, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT and DOUBLE. VARIANT is 
absent, and the converter switches in {{RawFormatDeserializationSchema}} and 
{{RawFormatSerializationSchema}} both end in {{UnsupportedOperationException}} 
for anything off that list.

Two consequences:
 # *Reading.* A schemaless JSON topic must be declared as {{{}STRING{}}}, and 
every query has to call {{PARSE_JSON}} on it. The schema then describes the 
column as text, which is not what it holds. Catalogs, {{DESCRIBE}} and 
downstream consumers all see the wrong type, and the parse is repeated per 
query instead of happening once at the source.
 # *Writing.* A VARIANT column cannot be written to a schemaless sink at all. 
{{INSERT INTO raw_sink SELECT v FROM ...}} is rejected, so the value has to be 
flattened with {{JSON_STRING}} first. Every other format that gains VARIANT 
support will accept it directly.

*Proposed change*

Treat a {{raw}} VARIANT column as "the bytes are a JSON document", decoded with 
the existing {{raw.charset}} option. Read is {{PARSE_JSON}} over the decoded 
bytes. Write is {{Variant#toJson}} encoded back with the same charset.
{code:java}
// RawFormatFactory#supportedTypes
LogicalTypeRoot.VARIANT,
{code}
{code:java}
// RawFormatDeserializationSchema#createConverter
case VARIANT:
    // new String(data, charset), then
    // BinaryVariantInternalBuilder.parseJson(json, false)
    return createVariantConverter(charsetName);

// RawFormatDeserializationSchema#createDataLengthValidator
case VARIANT:              // variable length, as CHAR / VARBINARY / RAW
    return data -> {};
{code}
{code:java}
// RawFormatSerializationSchema#createConverter
case VARIANT:
    return row -> row.getVariant(0).toJson().getBytes(charset);
{code}
The deserializer returns {{{}org.apache.flink.types.variant.Variant{}}}, which 
is what {{RowData#getVariant}} already contracts for, so no new data structure 
is needed. Docs need one row in the type table of 
{{docs/content/docs/connectors/table/formats/raw.md}} and its {{content.zh}} 
counterpart.

*Why JSON text and not the binary variant encoding*

A variant value is two byte arrays, {{value}} and {{{}metadata{}}}. A {{raw}} 
column carries exactly one {{{}byte[]{}}}. Interpreting the bytes as the 
Parquet/Spark variant binary encoding would mean inventing a framing to pack 
both blobs into one array, which is precisely the kind of wire convention the 
{{raw}} format exists to avoid. JSON text is the only interpretation that needs 
no new convention. Formats that can carry two fields, such as Avro or Protobuf, 
are the right home for the binary encoding.

*Design notes*
 * *Duplicate keys.* Parse with {{{}allowDuplicateKeys = false{}}}, matching 
the single-argument {{PARSE_JSON}} default. A format option could expose this 
later if asked for.
 * *Malformed input.* Throw, consistent with how the fixed-width types reject 
wrong-length input. The {{raw}} format has no {{ignore-parse-errors}} option, 
and adding one is a separate discussion.
 * *Options.* {{raw.charset}} applies. {{raw.endianness}} is meaningless for 
VARIANT, exactly as it is for CHAR and VARBINARY.
 * *Round-trip is value-lossless, not byte-lossless.* {{PARSE_JSON}} followed 
by {{toJson}} normalizes whitespace, key order and number spelling, so output 
bytes need not equal input bytes. Tests must assert on the parsed value, not on 
the bytes.
 * *Line delimiter.* The {{raw.line-delimiter}} option from FLINK-39401 
composes for free, which makes newline-delimited JSON readable as a stream of 
VARIANT rows.
 * *Key formats.* A VARIANT column has no canonical byte form, so it is a poor 
choice for a key format. No special handling is proposed.
 * *Efficiency.* The direct implementation materializes a {{String}} before 
parsing, because {{BinaryVariantInternalBuilder#parseJson}} accepts nothing 
else. FLINK-40218 removes that copy by letting a caller drive the builder from 
its own token source. This ticket is a direct beneficiary and should adopt it 
once available.

*Compatibility*

Purely additive. No existing type changes behavior. The only new failure mode 
appears where the plan previously failed with a {{{}ValidationException{}}}.

*Verifying this change*
 * {{{}RawFormatFactoryTest{}}}: a VARIANT column is accepted, unsupported 
types are still rejected, multiple physical columns are still rejected.
 * {{{}RawFormatSerDeSchemaTest{}}}: round-trip an object, an array and each 
JSON scalar form. A null message yields a null field. A non-UTF-8 charset 
works. Malformed JSON raises a deserialization error. Serializing a VARIANT 
produced by {{PARSE_JSON}} yields bytes that parse back to an equal value.
 * {{{}RawFormatLineDelimiterTest{}}}: a newline-delimited JSON message 
produces one VARIANT row per line.
```

  was:
*relates to* FLINK-37922, *relates to* FLINK-40218

 

The \{{raw}} format is the only way to read a topic or a file that carries no 
schema. VARIANT is the type Flink added for semi-structured data. The two 
cannot be combined today.

{code:sql}
CREATE TABLE events (payload VARIANT) WITH ('format' = 'raw', ...);
{code}
fails at plan time with:
{noformat}
org.apache.flink.table.api.ValidationException:
  The 'raw' format doesn't supports 'VARIANT' as column type.
{noformat}
{\{RawFormatFactory#supportedTypes}} lists CHAR, VARCHAR, BINARY, VARBINARY, 
RAW, BOOLEAN, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT and DOUBLE. VARIANT is 
absent, and the converter switches in \{{RawFormatDeserializationSchema}} and 
\{{RawFormatSerializationSchema}} both end in 
\{{UnsupportedOperationException}} for anything off that list.

Two consequences:
 # *Reading.* A schemaless JSON topic must be declared as \{{STRING}}, and 
every query has to call \{{PARSE_JSON}} on it. The schema then describes the 
column as text, which is not what it holds. Catalogs, \{{DESCRIBE}} and 
downstream consumers all see the wrong type, and the parse is repeated per 
query instead of happening once at the source.
 # *Writing.* A VARIANT column cannot be written to a schemaless sink at all. 
\{{INSERT INTO raw_sink SELECT v FROM ...}} is rejected, so the value has to be 
flattened with \{{JSON_STRING}} first. Every other format that gains VARIANT 
support will accept it directly.

*Proposed change*

Treat a \{{raw}} VARIANT column as "the bytes are a JSON document", decoded 
with the existing \{{raw.charset}} option. Read is \{{PARSE_JSON}} over the 
decoded bytes. Write is \{{Variant#toJson}} encoded back with the same charset.

{code:java}
// RawFormatFactory#supportedTypes
LogicalTypeRoot.VARIANT,
{code}

{code:java}
// RawFormatDeserializationSchema#createConverter
case VARIANT:
    // new String(data, charset), then
    // BinaryVariantInternalBuilder.parseJson(json, false)
    return createVariantConverter(charsetName);

// RawFormatDeserializationSchema#createDataLengthValidator
case VARIANT:              // variable length, as CHAR / VARBINARY / RAW
    return data -> {};
{code}

{code:java}
// RawFormatSerializationSchema#createConverter
case VARIANT:
    return row -> row.getVariant(0).toJson().getBytes(charset);
{code}

The deserializer returns \{{org.apache.flink.types.variant.Variant}}, which is 
what \{{RowData#getVariant}} already contracts for, so no new data structure is 
needed. Docs need one row in the type table of 
\{{docs/content/docs/connectors/table/formats/raw.md}} and its \{{content.zh}} 
counterpart.

*Why JSON text and not the binary variant encoding*

A variant value is two byte arrays, \{{value}} and \{{metadata}}. A \{{raw}} 
column carries exactly one \{{byte[]}}. Interpreting the bytes as the 
Parquet/Spark variant binary encoding would mean inventing a framing to pack 
both blobs into one array, which is precisely the kind of wire convention the 
\{{raw}} format exists to avoid. JSON text is the only interpretation that 
needs no new convention. Formats that can carry two fields, such as Avro or 
Protobuf, are the right home for the binary encoding.

*Design notes*
 * *Duplicate keys.* Parse with \{{allowDuplicateKeys = false}}, matching the 
single-argument \{{PARSE_JSON}} default. A format option could expose this 
later if asked for.
 * *Malformed input.* Throw, consistent with how the fixed-width types reject 
wrong-length input. The \{{raw}} format has no \{{ignore-parse-errors}} option, 
and adding one is a separate discussion.
 * *Options.* \{{raw.charset}} applies. \{{raw.endianness}} is meaningless for 
VARIANT, exactly as it is for CHAR and VARBINARY.
 * *Round-trip is value-lossless, not byte-lossless.* \{{PARSE_JSON}} followed 
by \{{toJson}} normalizes whitespace, key order and number spelling, so output 
bytes need not equal input bytes. Tests must assert on the parsed value, not on 
the bytes.
 * *Line delimiter.* The \{{raw.line-delimiter}} option from FLINK-39401 
composes for free, which makes newline-delimited JSON readable as a stream of 
VARIANT rows.
 * *Key formats.* A VARIANT column has no canonical byte form, so it is a poor 
choice for a key format. No special handling is proposed.
 * *Efficiency.* The direct implementation materializes a \{{String}} before 
parsing, because \{{BinaryVariantInternalBuilder#parseJson}} accepts nothing 
else. FLINK-40218 removes that copy by letting a caller drive the builder from 
its own token source. This ticket is a direct beneficiary and should adopt it 
once available.

*Compatibility*

Purely additive. No existing type changes behavior. The only new failure mode 
appears where the plan previously failed with a \{{ValidationException}}.

*Verifying this change*
 * \{{RawFormatFactoryTest}}: a VARIANT column is accepted, unsupported types 
are still rejected, multiple physical columns are still rejected.
 * \{{RawFormatSerDeSchemaTest}}: round-trip an object, an array and each JSON 
scalar form. A null message yields a null field. A non-UTF-8 charset works. 
Malformed JSON raises a deserialization error. Serializing a VARIANT produced 
by \{{PARSE_JSON}} yields bytes that parse back to an equal value.
 * \{{RawFormatLineDelimiterTest}}: a newline-delimited JSON message produces 
one VARIANT row per line.
```


> Support VARIANT as column type in the 'raw' format
> --------------------------------------------------
>
>                 Key: FLINK-40241
>                 URL: https://issues.apache.org/jira/browse/FLINK-40241
>             Project: Flink
>          Issue Type: Sub-task
>          Components: Formats (JSON, Avro, Parquet, ORC, SequenceFile)
>            Reporter: Ramin Gharib
>            Assignee: Ramin Gharib
>            Priority: Major
>
> *part to* FLINK-37922,  *relates to* FLINK-40218 , FLINK-39605, FLINK-39401
> The {{raw}} format is the only way to read a topic or a file that carries no 
> schema. VARIANT is the type Flink added for semi-structured data. The two 
> cannot be combined today.
> {code:sql}
> CREATE TABLE events (payload VARIANT) WITH ('format' = 'raw', ...);
> {code}
> fails at plan time with:
> {noformat}
> org.apache.flink.table.api.ValidationException:
>   The 'raw' format doesn't supports 'VARIANT' as column type.
> {noformat}
> {{RawFormatFactory#supportedTypes}} lists CHAR, VARCHAR, BINARY, VARBINARY, 
> RAW, BOOLEAN, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT and DOUBLE. VARIANT 
> is absent, and the converter switches in {{RawFormatDeserializationSchema}} 
> and {{RawFormatSerializationSchema}} both end in 
> {{UnsupportedOperationException}} for anything off that list.
> Two consequences:
>  # *Reading.* A schemaless JSON topic must be declared as {{{}STRING{}}}, and 
> every query has to call {{PARSE_JSON}} on it. The schema then describes the 
> column as text, which is not what it holds. Catalogs, {{DESCRIBE}} and 
> downstream consumers all see the wrong type, and the parse is repeated per 
> query instead of happening once at the source.
>  # *Writing.* A VARIANT column cannot be written to a schemaless sink at all. 
> {{INSERT INTO raw_sink SELECT v FROM ...}} is rejected, so the value has to 
> be flattened with {{JSON_STRING}} first. Every other format that gains 
> VARIANT support will accept it directly.
> *Proposed change*
> Treat a {{raw}} VARIANT column as "the bytes are a JSON document", decoded 
> with the existing {{raw.charset}} option. Read is {{PARSE_JSON}} over the 
> decoded bytes. Write is {{Variant#toJson}} encoded back with the same charset.
> {code:java}
> // RawFormatFactory#supportedTypes
> LogicalTypeRoot.VARIANT,
> {code}
> {code:java}
> // RawFormatDeserializationSchema#createConverter
> case VARIANT:
>     // new String(data, charset), then
>     // BinaryVariantInternalBuilder.parseJson(json, false)
>     return createVariantConverter(charsetName);
> // RawFormatDeserializationSchema#createDataLengthValidator
> case VARIANT:              // variable length, as CHAR / VARBINARY / RAW
>     return data -> {};
> {code}
> {code:java}
> // RawFormatSerializationSchema#createConverter
> case VARIANT:
>     return row -> row.getVariant(0).toJson().getBytes(charset);
> {code}
> The deserializer returns {{{}org.apache.flink.types.variant.Variant{}}}, 
> which is what {{RowData#getVariant}} already contracts for, so no new data 
> structure is needed. Docs need one row in the type table of 
> {{docs/content/docs/connectors/table/formats/raw.md}} and its {{content.zh}} 
> counterpart.
> *Why JSON text and not the binary variant encoding*
> A variant value is two byte arrays, {{value}} and {{{}metadata{}}}. A {{raw}} 
> column carries exactly one {{{}byte[]{}}}. Interpreting the bytes as the 
> Parquet/Spark variant binary encoding would mean inventing a framing to pack 
> both blobs into one array, which is precisely the kind of wire convention the 
> {{raw}} format exists to avoid. JSON text is the only interpretation that 
> needs no new convention. Formats that can carry two fields, such as Avro or 
> Protobuf, are the right home for the binary encoding.
> *Design notes*
>  * *Duplicate keys.* Parse with {{{}allowDuplicateKeys = false{}}}, matching 
> the single-argument {{PARSE_JSON}} default. A format option could expose this 
> later if asked for.
>  * *Malformed input.* Throw, consistent with how the fixed-width types reject 
> wrong-length input. The {{raw}} format has no {{ignore-parse-errors}} option, 
> and adding one is a separate discussion.
>  * *Options.* {{raw.charset}} applies. {{raw.endianness}} is meaningless for 
> VARIANT, exactly as it is for CHAR and VARBINARY.
>  * *Round-trip is value-lossless, not byte-lossless.* {{PARSE_JSON}} followed 
> by {{toJson}} normalizes whitespace, key order and number spelling, so output 
> bytes need not equal input bytes. Tests must assert on the parsed value, not 
> on the bytes.
>  * *Line delimiter.* The {{raw.line-delimiter}} option from FLINK-39401 
> composes for free, which makes newline-delimited JSON readable as a stream of 
> VARIANT rows.
>  * *Key formats.* A VARIANT column has no canonical byte form, so it is a 
> poor choice for a key format. No special handling is proposed.
>  * *Efficiency.* The direct implementation materializes a {{String}} before 
> parsing, because {{BinaryVariantInternalBuilder#parseJson}} accepts nothing 
> else. FLINK-40218 removes that copy by letting a caller drive the builder 
> from its own token source. This ticket is a direct beneficiary and should 
> adopt it once available.
> *Compatibility*
> Purely additive. No existing type changes behavior. The only new failure mode 
> appears where the plan previously failed with a {{{}ValidationException{}}}.
> *Verifying this change*
>  * {{{}RawFormatFactoryTest{}}}: a VARIANT column is accepted, unsupported 
> types are still rejected, multiple physical columns are still rejected.
>  * {{{}RawFormatSerDeSchemaTest{}}}: round-trip an object, an array and each 
> JSON scalar form. A null message yields a null field. A non-UTF-8 charset 
> works. Malformed JSON raises a deserialization error. Serializing a VARIANT 
> produced by {{PARSE_JSON}} yields bytes that parse back to an equal value.
>  * {{{}RawFormatLineDelimiterTest{}}}: a newline-delimited JSON message 
> produces one VARIANT row per line.
> ```



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to