This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 8df195315e1 [opt](paimon) support read paimon variant type using jni
reader (#66334)
8df195315e1 is described below
commit 8df195315e1f26d5947ed14806087c5ab4cb38ca
Author: zhangstar333 <[email protected]>
AuthorDate: Tue Aug 4 21:22:57 2026 +0800
[opt](paimon) support read paimon variant type using jni reader (#66334)
### What problem does this PR solve?
Problem Summary:
support read paimon variant type using jni reader
```
mysql> desc test_paimon_variant.db1.variant_smoke;
+---------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------+------+---------+-------+
| Field | Type
| Null | Key |
Default | Extra |
+---------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------+------+---------+-------+
| id | bigint
| Yes | true | NULL
| |
| payload | variant<PROPERTIES ("variant_max_subcolumns_count" =
"0","variant_enable_typed_paths_to_sparse" =
"false","variant_max_sparse_column_statistics_size" =
"10000","variant_sparse_hash_shard_count" = "1")> | Yes | true | NULL |
|
+---------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+------+------+---------+-------+
2 rows in set (0.01 sec)
mysql> select * from test_paimon_variant.db1.variant_smoke;
+------+-----------------------------------------------------------------------------------------------------------------------------------------+
| id | payload
|
+------+-----------------------------------------------------------------------------------------------------------------------------------------+
| 2 |
{"active":false,"age":30,"extra":{"levels":[1,2,3]},"name":"bob","profile":{"city":"shanghai"},"tags":["doris"]}
|
| 1 |
{"active":true,"age":18,"missing":null,"name":"alice","profile":{"city":"beijing","zip":100000},"score":98.5,"tags":["flink","paimon"]}
|
| 3 | [1,"mixed",false,null,{"k":"v"}]
|
+------+-----------------------------------------------------------------------------------------------------------------------------------------+
3 rows in set (0.13 sec)
```
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [x] Regression test
- [ ] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
be/src/format/jni/jni_data_bridge.cpp | 29 +++
be/src/format/jni/jni_data_bridge.h | 3 +
be/test/format/jni/jni_data_bridge_test.cpp | 89 +++++++++
.../create_preinstalled_scripts/paimon/run13.sql | 17 ++
.../apache/doris/common/jni/vec/ColumnType.java | 12 ++
.../apache/doris/common/jni/vec/ColumnValue.java | 8 +
.../apache/doris/common/jni/vec/VectorColumn.java | 31 +++
.../doris/common/jni/vec/VectorColumnVariant.java | 214 +++++++++++++++++++++
.../common/jni/vec/VectorColumnVariantTest.java | 128 ++++++++++++
.../org/apache/doris/paimon/PaimonColumnValue.java | 15 ++
.../org/apache/doris/paimon/PaimonTypeUtils.java | 6 +
.../apache/doris/paimon/PaimonColumnValueTest.java | 15 ++
.../java/org/apache/doris/catalog/VariantType.java | 6 +
.../apache/doris/datasource/paimon/PaimonUtil.java | 5 +-
.../doris/datasource/paimon/PaimonUtilTest.java | 10 +
.../paimon/test_paimon_catalog_variant.out | 26 +++
.../paimon/test_paimon_catalog_variant.groovy | 91 +++++++++
17 files changed, 703 insertions(+), 2 deletions(-)
diff --git a/be/src/format/jni/jni_data_bridge.cpp
b/be/src/format/jni/jni_data_bridge.cpp
index 9dc935e0b62..adcf11e196e 100644
--- a/be/src/format/jni/jni_data_bridge.cpp
+++ b/be/src/format/jni/jni_data_bridge.cpp
@@ -29,6 +29,7 @@
#include "core/column/column_string.h"
#include "core/column/column_struct.h"
#include "core/column/column_varbinary.h"
+#include "core/column/variant_v2/column_variant_v2.h"
#include "core/data_type/data_type_array.h"
#include "core/data_type/data_type_map.h"
#include "core/data_type/data_type_nullable.h"
@@ -148,6 +149,9 @@ Status JniDataBridge::fill_column(TableMetaAddress&
address, ColumnPtr& doris_co
case PrimitiveType::TYPE_VARBINARY:
status = _fill_varbinary_column(address, data_column, num_rows);
break;
+ case PrimitiveType::TYPE_VARIANT:
+ status = _fill_variant_v2_column(address, data_column, num_rows);
+ break;
default:
status = Status::InvalidArgument("Unsupported type {} in jni scanner",
data_type->get_name());
@@ -179,6 +183,27 @@ Status
JniDataBridge::_fill_varbinary_column(TableMetaAddress& address,
return Status::OK();
}
+Status JniDataBridge::_fill_variant_v2_column(TableMetaAddress& address,
+ MutableColumnPtr& doris_column,
size_t num_rows) {
+ // VectorColumnVariant publishes these EncodedDataView fields immediately
after the null map.
+ const auto metadata_count =
static_cast<size_t>(address.next_meta_as_long());
+ const auto* metadata_offsets = reinterpret_cast<const
uint32_t*>(address.next_meta_as_ptr());
+ const auto* metadata_bytes = reinterpret_cast<const
char*>(address.next_meta_as_ptr());
+ const auto* metadata_ids = reinterpret_cast<const
uint32_t*>(address.next_meta_as_ptr());
+ const auto* value_offsets = reinterpret_cast<const
uint32_t*>(address.next_meta_as_ptr());
+ const auto* value_bytes = reinterpret_cast<const
char*>(address.next_meta_as_ptr());
+
+ auto& variant_column = assert_cast<ColumnVariantV2&>(*doris_column);
+ variant_column.insert_encoded_rows({
+ .metadata_bytes = {metadata_bytes,
metadata_offsets[metadata_count]},
+ .metadata_offsets = {metadata_offsets, metadata_count + 1},
+ .meta_ids = {metadata_ids, num_rows},
+ .value_bytes = {value_bytes, value_offsets[num_rows]},
+ .value_offsets = {value_offsets, num_rows + 1},
+ });
+ return Status::OK();
+}
+
Status JniDataBridge::_fill_string_column(TableMetaAddress& address,
MutableColumnPtr& doris_column,
size_t num_rows) {
auto& string_col = static_cast<ColumnString&>(*doris_column);
@@ -358,6 +383,8 @@ std::string JniDataBridge::get_jni_type(const DataTypePtr&
data_type) {
}
case TYPE_VARBINARY:
return "varbinary";
+ case TYPE_VARIANT:
+ return "variant";
// bitmap, hll, quantile_state, jsonb are transferred as strings via JNI
case TYPE_BITMAP:
[[fallthrough]];
@@ -433,6 +460,8 @@ std::string
JniDataBridge::get_jni_type_with_different_string(const DataTypePtr&
<< assert_cast<const
DataTypeVarbinary*>(remove_nullable(data_type).get())->len()
<< ")";
return buffer.str();
+ case TYPE_VARIANT:
+ return "variant";
case TYPE_DECIMALV2: {
buffer << "decimalv2(" << DecimalV2Value::PRECISION << "," <<
DecimalV2Value::SCALE << ")";
return buffer.str();
diff --git a/be/src/format/jni/jni_data_bridge.h
b/be/src/format/jni/jni_data_bridge.h
index e037ffec3d4..267d0a1711c 100644
--- a/be/src/format/jni/jni_data_bridge.h
+++ b/be/src/format/jni/jni_data_bridge.h
@@ -154,6 +154,9 @@ private:
static Status _fill_varbinary_column(TableMetaAddress& address,
MutableColumnPtr& doris_column,
size_t num_rows);
+ static Status _fill_variant_v2_column(TableMetaAddress& address,
MutableColumnPtr& doris_column,
+ size_t num_rows);
+
static Status _fill_array_column(TableMetaAddress& address,
MutableColumnPtr& doris_column,
const DataTypePtr& data_type, size_t
num_rows);
diff --git a/be/test/format/jni/jni_data_bridge_test.cpp
b/be/test/format/jni/jni_data_bridge_test.cpp
new file mode 100644
index 00000000000..d6418c0b97e
--- /dev/null
+++ b/be/test/format/jni/jni_data_bridge_test.cpp
@@ -0,0 +1,89 @@
+// 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.
+
+#include "format/jni/jni_data_bridge.h"
+
+#include <gtest/gtest.h>
+
+#include <array>
+#include <memory>
+
+#include "core/column/column_nullable.h"
+#include "core/column/variant_v2/column_variant_v2.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_variant_v2.h"
+
+namespace doris {
+namespace {
+
+template <typename T>
+long address_of(T* pointer) {
+ return reinterpret_cast<long>(pointer);
+}
+
+TEST(JniDataBridgeTest, VariantTypeUsesVariantJniName) {
+ auto variant_type = std::make_shared<DataTypeVariantV2>();
+ auto nullable_variant_type = make_nullable(variant_type);
+
+ EXPECT_EQ(JniDataBridge::get_jni_type(variant_type), "variant");
+ EXPECT_EQ(JniDataBridge::get_jni_type(nullable_variant_type), "variant");
+ EXPECT_EQ(JniDataBridge::get_jni_type_with_different_string(variant_type),
"variant");
+
EXPECT_EQ(JniDataBridge::get_jni_type_with_different_string(nullable_variant_type),
"variant");
+}
+
+TEST(JniDataBridgeTest, FillNullableVariantV2FromEncodedPaimonLayout) {
+ std::array<bool, 3> null_map {false, true, false};
+ // Paimon metadata contains an empty dictionary and a one-key dictionary
for "a".
+ std::array<uint32_t, 3> metadata_offsets {0, 3, 8};
+ std::array<char, 8> metadata_bytes {char {0x01}, char {0x00}, char {0x00},
char {0x01},
+ char {0x01}, char {0x00}, char {0x01},
'a'};
+ std::array<uint32_t, 3> metadata_ids {0, 0, 1};
+ // Primitive values are true, Variant null (the SQL NULL placeholder), and
false.
+ std::array<uint32_t, 4> value_offsets {0, 1, 2, 3};
+ std::array<char, 3> value_bytes {char {0x04}, char {0x00}, char {0x08}};
+
+ std::array<long, 7> meta {
+ address_of(null_map.data()), 2,
+ address_of(metadata_offsets.data()),
address_of(metadata_bytes.data()),
+ address_of(metadata_ids.data()),
address_of(value_offsets.data()),
+ address_of(value_bytes.data()),
+ };
+ JniDataBridge::TableMetaAddress meta_address(address_of(meta.data()));
+
+ DataTypePtr data_type =
make_nullable(std::make_shared<DataTypeVariantV2>());
+ ColumnPtr column = data_type->create_column();
+ const Status status = JniDataBridge::fill_column(meta_address, column,
data_type, 3);
+ ASSERT_TRUE(status.ok()) << status;
+
+ const auto& nullable_column = assert_cast<const ColumnNullable&>(*column);
+ EXPECT_EQ(nullable_column.get_null_map_data(), NullMap({0, 1, 0}));
+
+ const auto& variant_column =
+ assert_cast<const
ColumnVariantV2&>(nullable_column.get_nested_column());
+ ASSERT_EQ(variant_column.size(), 3);
+ ASSERT_FALSE(variant_column.is_typed());
+ const auto view = variant_column.read_view();
+ EXPECT_EQ(view.metadata_count(), 2);
+ EXPECT_EQ(view.metadata_id_at(0), view.metadata_id_at(1));
+ EXPECT_NE(view.metadata_id_at(0), view.metadata_id_at(2));
+ EXPECT_TRUE(variant_column.get_value_ref(0).get_bool());
+ EXPECT_TRUE(variant_column.get_value_ref(1).is_null());
+ EXPECT_FALSE(variant_column.get_value_ref(2).get_bool());
+}
+
+} // namespace
+} // namespace doris
diff --git
a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql
new file mode 100644
index 00000000000..51ad9973ce9
--- /dev/null
+++
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/paimon/run13.sql
@@ -0,0 +1,17 @@
+use paimon;
+create database if not exists test_paimon_spark;
+use test_paimon_spark;
+
+drop table if exists variant_smoke;
+create table variant_smoke (
+ id BIGINT,
+ payload VARIANT
+) using paimon
+tblproperties (
+ 'file.format' = 'parquet'
+);
+
+insert into variant_smoke values
+ (1,
parse_json('{"name":"alice","age":18,"active":true,"score":98.5,"tags":["flink","paimon"],"profile":{"city":"beijing","zip":100000},"missing":null}')),
+ (2,
parse_json('{"name":"bob","age":30,"active":false,"tags":["doris"],"profile":{"city":"shanghai"},"extra":{"levels":[1,2,3]}}')),
+ (3, parse_json('[1,"mixed",false,null,{"k":"v"}]'));
diff --git
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnType.java
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnType.java
index 983681d24dc..00411e85312 100644
---
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnType.java
+++
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnType.java
@@ -63,6 +63,7 @@ public class ColumnType {
IPV6(16),
STRING(-1),
VARBINARY(-1),
+ VARIANT(-1),
ARRAY(-1),
MAP(-1),
STRUCT(-1);
@@ -157,6 +158,10 @@ public class ColumnType {
return type == Type.BINARY || type == Type.VARBINARY;
}
+ public boolean isVariantType() {
+ return type == Type.VARIANT;
+ }
+
public boolean isComplexType() {
return type == Type.ARRAY || type == Type.MAP || type == Type.STRUCT;
}
@@ -250,6 +255,10 @@ public class ColumnType {
case VARCHAR:
// [const | nullMap | offsets | data ]
return 4;
+ case VARIANT:
+ // [const | nullMap | metadata count | metadata offsets |
metadata bytes
+ // | metadata ids | value offsets | value bytes]
+ return 8;
default:
// [const | nullMap | data]
return 3;
@@ -344,6 +353,9 @@ public class ColumnType {
case "varbinary":
type = Type.VARBINARY;
break;
+ case "variant":
+ type = Type.VARIANT;
+ break;
default:
if (lowerCaseType.startsWith("timestamptz")) {
type = Type.TIMESTAMPTZ;
diff --git
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnValue.java
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnValue.java
index 8911a19f0bd..36c0be38636 100644
---
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnValue.java
+++
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/ColumnValue.java
@@ -77,6 +77,14 @@ public interface ColumnValue {
byte[] getBytes();
+ default byte[] getVariantMetadata() {
+ throw new UnsupportedOperationException("Variant metadata is not
available");
+ }
+
+ default byte[] getVariantValue() {
+ throw new UnsupportedOperationException("Variant value is not
available");
+ }
+
void unpackArray(List<ColumnValue> values);
void unpackMap(List<ColumnValue> keys, List<ColumnValue> values);
diff --git
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumn.java
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumn.java
index 36d2329da89..010bf7eba8f 100644
---
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumn.java
+++
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumn.java
@@ -75,6 +75,7 @@ public class VectorColumn {
// For nested column type: String / Array/ Map / Struct
private VectorColumn[] childColumns;
+ private VectorColumnVariant variantColumn;
// For struct, only support to read all fields in struct now
// todo: support pruned struct fields
@@ -120,6 +121,8 @@ public class VectorColumn {
childColumns = new VectorColumn[1];
childColumns[0] = new VectorColumn(new ColumnType("#stringBytes",
Type.BYTE),
capacity * DEFAULT_STRING_LENGTH);
+ } else if (columnType.isVariantType()) {
+ variantColumn = new VectorColumnVariant();
}
reserveCapacity(capacity);
@@ -273,6 +276,10 @@ public class VectorColumn {
}
childColumns = null;
}
+ if (variantColumn != null) {
+ variantColumn.close();
+ variantColumn = null;
+ }
if (nullMap != 0) {
OffHeap.freeMemory(nullMap);
@@ -354,6 +361,8 @@ public class VectorColumn {
this.offsets = OffHeap.reallocateMemory(offsets, oldOffsetSize,
newOffsetSize);
} else if (columnType.isVarbinaryType()) {
this.data = OffHeap.reallocateMemory(data, oldCapacity * 16L,
newCapacity * 16L);
+ } else if (columnType.isVariantType()) {
+ variantColumn.reserveRows(newCapacity);
} else if (!columnType.isStruct()) {
throw new RuntimeException("Unhandled type: " +
columnType.getName());
}
@@ -370,6 +379,9 @@ public class VectorColumn {
c.reset();
}
}
+ if (variantColumn != null) {
+ variantColumn.reset();
+ }
appendIndex = 0;
if (numNulls > 0) {
putNotNulls(0, capacity);
@@ -462,6 +474,8 @@ public class VectorColumn {
case BINARY:
case VARBINARY:
return appendVarbinary(new byte[0]);
+ case VARIANT:
+ return appendVariantNull();
default:
throw new RuntimeException("Unknown type value: " + typeValue);
}
@@ -1530,6 +1544,17 @@ public class VectorColumn {
return appendIndex++;
}
+ public int appendVariant(byte[] metadata, byte[] value) {
+ reserve(appendIndex + 1);
+ variantColumn.append(metadata, value);
+ return appendIndex++;
+ }
+
+ private int appendVariantNull() {
+ variantColumn.appendNull();
+ return appendIndex++;
+ }
+
public void appendVarbinary(byte[][] batch, boolean isNullable) {
if (!isNullable) {
checkNullable(batch, batch.length);
@@ -1619,6 +1644,9 @@ public class VectorColumn {
for (VectorColumn c : childColumns) {
c.updateMeta(meta);
}
+ } else if (columnType.isVariantType()) {
+ meta.appendLong(nullMap);
+ variantColumn.updateMeta(meta);
} else {
meta.appendLong(nullMap);
meta.appendLong(data);
@@ -1903,6 +1931,9 @@ public class VectorColumn {
case VARBINARY:
appendVarbinary(o.getBytes());
break;
+ case VARIANT:
+ appendVariant(o.getVariantMetadata(), o.getVariantValue());
+ break;
case ARRAY: {
List<ColumnValue> values = new ArrayList<>();
o.unpackArray(values);
diff --git
a/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumnVariant.java
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumnVariant.java
new file mode 100644
index 00000000000..91e66fe2e9f
--- /dev/null
+++
b/fe/be-java-extensions/java-common/src/main/java/org/apache/doris/common/jni/vec/VectorColumnVariant.java
@@ -0,0 +1,214 @@
+// 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.doris.common.jni.vec;
+
+import org.apache.doris.common.jni.utils.OffHeap;
+
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Objects;
+
+/**
+ * Off-heap encoded Variant storage used by {@link VectorColumn}.
+ *
+ * <p>The layout mirrors ColumnVariantV2::EncodedDataView: a deduplicated
metadata dictionary,
+ * one metadata id per row, and one encoded value per row.</p>
+ */
+final class VectorColumnVariant {
+ private static final int MAX_CAPACITY = Integer.MAX_VALUE - 15;
+ private static final byte[] EMPTY_METADATA = new byte[] {1, 0, 0};
+ private static final byte[] NULL_VALUE = new byte[] {0};
+
+ private final Map<ByteArrayKey, Integer> metadataIds = new HashMap<>();
+
+ private long metadataOffsets;
+ private long metadataBytes;
+ private long rowMetadataIds;
+ private long valueOffsets;
+ private long valueBytes;
+
+ private int rowCapacity;
+ private int metadataBytesCapacity;
+ private int valueBytesCapacity;
+
+ private int metadataCount;
+ private int metadataBytesSize;
+ private int rowCount;
+ private int valueBytesSize;
+
+ void reserveRows(int requiredCapacity) {
+ if (requiredCapacity <= rowCapacity) {
+ return;
+ }
+ int newCapacity = growCapacity(rowCapacity, requiredCapacity);
+ long oldOffsetsSize = rowCapacity == 0 ? 0 : (rowCapacity + 1L) *
Integer.BYTES;
+ metadataOffsets = OffHeap.reallocateMemory(
+ metadataOffsets, oldOffsetsSize,
+ (newCapacity + 1L) * Integer.BYTES);
+ rowMetadataIds = OffHeap.reallocateMemory(
+ rowMetadataIds, (long) rowCapacity * Integer.BYTES,
+ (long) newCapacity * Integer.BYTES);
+ valueOffsets = OffHeap.reallocateMemory(
+ valueOffsets, oldOffsetsSize,
+ (newCapacity + 1L) * Integer.BYTES);
+ if (rowCapacity == 0) {
+ OffHeap.putInt(null, metadataOffsets, 0);
+ OffHeap.putInt(null, valueOffsets, 0);
+ }
+ rowCapacity = newCapacity;
+ }
+
+ void append(byte[] metadata, byte[] value) {
+ Objects.requireNonNull(metadata, "Variant metadata cannot be null");
+ Objects.requireNonNull(value, "Variant value cannot be null");
+ reserveRows(rowCount + 1);
+
+ Integer metadataId = metadataIds.get(new ByteArrayKey(metadata));
+ if (metadataId == null) {
+ metadataId = appendMetadata(metadata);
+ }
+ OffHeap.putInt(null, rowMetadataIds + (long) rowCount * Integer.BYTES,
metadataId);
+
+ int requiredValueBytes = checkedSize("value", valueBytesSize,
value.length);
+ reserveValueBytes(requiredValueBytes);
+ OffHeap.copyMemory(value, OffHeap.BYTE_ARRAY_OFFSET, null, valueBytes
+ valueBytesSize, value.length);
+ valueBytesSize = requiredValueBytes;
+ rowCount++;
+ OffHeap.putInt(null, valueOffsets + (long) rowCount * Integer.BYTES,
valueBytesSize);
+ }
+
+ void appendNull() {
+ append(EMPTY_METADATA, NULL_VALUE);
+ }
+
+ void updateMeta(VectorColumn meta) {
+ meta.appendLong(metadataCount);
+ meta.appendLong(metadataOffsets);
+ meta.appendLong(metadataBytes);
+ meta.appendLong(rowMetadataIds);
+ meta.appendLong(valueOffsets);
+ meta.appendLong(valueBytes);
+ }
+
+ void reset() {
+ metadataIds.clear();
+ metadataCount = 0;
+ metadataBytesSize = 0;
+ rowCount = 0;
+ valueBytesSize = 0;
+ if (rowCapacity > 0) {
+ OffHeap.putInt(null, metadataOffsets, 0);
+ OffHeap.putInt(null, valueOffsets, 0);
+ }
+ }
+
+ void close() {
+ free(metadataOffsets);
+ free(metadataBytes);
+ free(rowMetadataIds);
+ free(valueOffsets);
+ free(valueBytes);
+ metadataOffsets = 0;
+ metadataBytes = 0;
+ rowMetadataIds = 0;
+ valueOffsets = 0;
+ valueBytes = 0;
+ rowCapacity = 0;
+ metadataBytesCapacity = 0;
+ valueBytesCapacity = 0;
+ reset();
+ }
+
+ private int appendMetadata(byte[] metadata) {
+ int requiredMetadataBytes = checkedSize("metadata", metadataBytesSize,
metadata.length);
+ reserveMetadataBytes(requiredMetadataBytes);
+ OffHeap.copyMemory(
+ metadata, OffHeap.BYTE_ARRAY_OFFSET, null, metadataBytes +
metadataBytesSize, metadata.length);
+ metadataBytesSize = requiredMetadataBytes;
+ metadataCount++;
+ OffHeap.putInt(
+ null, metadataOffsets + (long) metadataCount * Integer.BYTES,
metadataBytesSize);
+
+ int metadataId = metadataCount - 1;
+ metadataIds.put(new ByteArrayKey(Arrays.copyOf(metadata,
metadata.length)), metadataId);
+ return metadataId;
+ }
+
+ private void reserveMetadataBytes(int requiredCapacity) {
+ if (requiredCapacity <= metadataBytesCapacity) {
+ return;
+ }
+ int newCapacity = growCapacity(metadataBytesCapacity,
requiredCapacity);
+ metadataBytes = OffHeap.reallocateMemory(metadataBytes,
metadataBytesCapacity, newCapacity);
+ metadataBytesCapacity = newCapacity;
+ }
+
+ private void reserveValueBytes(int requiredCapacity) {
+ if (requiredCapacity <= valueBytesCapacity) {
+ return;
+ }
+ int newCapacity = growCapacity(valueBytesCapacity, requiredCapacity);
+ valueBytes = OffHeap.reallocateMemory(valueBytes, valueBytesCapacity,
newCapacity);
+ valueBytesCapacity = newCapacity;
+ }
+
+ private static int checkedSize(String component, int currentSize, int
appendedSize) {
+ long requiredSize = (long) currentSize + appendedSize;
+ if (requiredSize > MAX_CAPACITY) {
+ throw new RuntimeException("Variant " + component + " buffer
exceeds the Java JNI size limit");
+ }
+ return (int) requiredSize;
+ }
+
+ private static int growCapacity(int currentCapacity, int requiredCapacity)
{
+ long doubledCapacity = Math.max(1L, currentCapacity * 2L);
+ int newCapacity = (int) Math.min(MAX_CAPACITY,
Math.max(doubledCapacity, requiredCapacity));
+ if (newCapacity < requiredCapacity) {
+ throw new RuntimeException("Cannot reserve enough bytes for
Variant JNI data");
+ }
+ return newCapacity;
+ }
+
+ private static void free(long address) {
+ if (address != 0) {
+ OffHeap.freeMemory(address);
+ }
+ }
+
+ private static final class ByteArrayKey {
+ private final byte[] bytes;
+ private final int hashCode;
+
+ private ByteArrayKey(byte[] bytes) {
+ this.bytes = bytes;
+ this.hashCode = Arrays.hashCode(bytes);
+ }
+
+ @Override
+ public boolean equals(Object other) {
+ return other instanceof ByteArrayKey
+ && Arrays.equals(bytes, ((ByteArrayKey) other).bytes);
+ }
+
+ @Override
+ public int hashCode() {
+ return hashCode;
+ }
+ }
+}
diff --git
a/fe/be-java-extensions/java-common/src/test/java/org/apache/doris/common/jni/vec/VectorColumnVariantTest.java
b/fe/be-java-extensions/java-common/src/test/java/org/apache/doris/common/jni/vec/VectorColumnVariantTest.java
new file mode 100644
index 00000000000..a95608e14a8
--- /dev/null
+++
b/fe/be-java-extensions/java-common/src/test/java/org/apache/doris/common/jni/vec/VectorColumnVariantTest.java
@@ -0,0 +1,128 @@
+// 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.doris.common.jni.vec;
+
+import org.apache.doris.common.jni.utils.OffHeap;
+
+import org.junit.Assert;
+import org.junit.BeforeClass;
+import org.junit.Test;
+
+public class VectorColumnVariantTest {
+ private static final byte[] EMPTY_METADATA = new byte[] {1, 0, 0};
+ private static final byte[] ONE_KEY_METADATA = new byte[] {1, 1, 0, 1,
'a'};
+
+ @BeforeClass
+ public static void setUpClass() {
+ OffHeap.setTesting();
+ }
+
+ @Test
+ public void testVariantTypeAndEncodedLayout() {
+ ColumnType variantType = ColumnType.parseType("v", "variant");
+ Assert.assertEquals(ColumnType.Type.VARIANT, variantType.getType());
+ Assert.assertEquals(8, variantType.metaSize());
+
+ VectorTable table = VectorTable.createWritableTable(
+ new ColumnType[] {variantType}, new String[] {"v"}, 2);
+ try {
+ VectorColumn column = table.getColumn(0);
+ column.appendVariant(EMPTY_METADATA, new byte[] {0});
+ column.appendVariant(EMPTY_METADATA.clone(), new byte[] {4});
+ column.appendVariant(ONE_KEY_METADATA, new byte[] {8});
+
+ long meta = table.getMetaAddress();
+ Assert.assertEquals(3L, OffHeap.getLong(null, meta));
+ Assert.assertArrayEquals(new boolean[] {false, false, false},
+ OffHeap.getBoolean(null, OffHeap.getLong(null, meta + 8),
3));
+ Assert.assertEquals(2L, OffHeap.getLong(null, meta + 16));
+
+ long metadataOffsets = OffHeap.getLong(null, meta + 24);
+ long metadataBytes = OffHeap.getLong(null, meta + 32);
+ long metadataIds = OffHeap.getLong(null, meta + 40);
+ long valueOffsets = OffHeap.getLong(null, meta + 48);
+ long valueBytes = OffHeap.getLong(null, meta + 56);
+
+ Assert.assertArrayEquals(new int[] {0, 3, 8},
+ OffHeap.getInt(null, metadataOffsets, 3));
+ Assert.assertArrayEquals(
+ new byte[] {1, 0, 0, 1, 1, 0, 1, 'a'},
+ OffHeap.getByte(null, metadataBytes, 8));
+ Assert.assertArrayEquals(new int[] {0, 0, 1},
+ OffHeap.getInt(null, metadataIds, 3));
+ Assert.assertArrayEquals(new int[] {0, 1, 2, 3},
+ OffHeap.getInt(null, valueOffsets, 4));
+ Assert.assertArrayEquals(new byte[] {0, 4, 8},
+ OffHeap.getByte(null, valueBytes, 3));
+ } finally {
+ table.close();
+ }
+ }
+
+ @Test
+ public void testSqlNullUsesValidVariantNullPlaceholder() {
+ ColumnType variantType = ColumnType.parseType("v", "variant");
+ VectorTable table = VectorTable.createWritableTable(
+ new ColumnType[] {variantType}, new String[] {"v"}, 1);
+ try {
+ VectorColumn column = table.getColumn(0);
+ column.appendVariant(EMPTY_METADATA, new byte[] {4});
+ column.appendNull(ColumnType.Type.VARIANT);
+
+ long meta = table.getMetaAddress();
+ long nullMap = OffHeap.getLong(null, meta + 8);
+ long metadataIds = OffHeap.getLong(null, meta + 40);
+ long valueOffsets = OffHeap.getLong(null, meta + 48);
+ long valueBytes = OffHeap.getLong(null, meta + 56);
+
+ Assert.assertArrayEquals(new boolean[] {false, true},
+ OffHeap.getBoolean(null, nullMap, 2));
+ Assert.assertEquals(1L, OffHeap.getLong(null, meta + 16));
+ Assert.assertArrayEquals(new int[] {0, 0},
+ OffHeap.getInt(null, metadataIds, 2));
+ Assert.assertArrayEquals(new int[] {0, 1, 2},
+ OffHeap.getInt(null, valueOffsets, 3));
+ Assert.assertArrayEquals(new byte[] {4, 0},
+ OffHeap.getByte(null, valueBytes, 2));
+ } finally {
+ table.close();
+ }
+ }
+
+ @Test
+ public void testResetRebuildsMetadataDictionary() {
+ ColumnType variantType = ColumnType.parseType("v", "variant");
+ VectorTable table = VectorTable.createWritableTable(
+ new ColumnType[] {variantType}, new String[] {"v"}, 1);
+ try {
+ table.getColumn(0).appendVariant(EMPTY_METADATA, new byte[] {0});
+ table.reset();
+ table.getColumn(0).appendVariant(ONE_KEY_METADATA, new byte[] {4});
+
+ long meta = table.getMetaAddress();
+ Assert.assertEquals(1L, OffHeap.getLong(null, meta));
+ Assert.assertEquals(1L, OffHeap.getLong(null, meta + 16));
+ Assert.assertArrayEquals(new int[] {0, 5},
+ OffHeap.getInt(null, OffHeap.getLong(null, meta + 24), 2));
+ Assert.assertArrayEquals(new int[] {0},
+ OffHeap.getInt(null, OffHeap.getLong(null, meta + 40), 1));
+ } finally {
+ table.close();
+ }
+ }
+}
diff --git
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java
index 339feb97c9b..d92983f5c0c 100644
---
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java
+++
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonColumnValue.java
@@ -25,6 +25,7 @@ import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.data.variant.Variant;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.LocalZonedTimestampType;
@@ -198,6 +199,20 @@ public class PaimonColumnValue implements ColumnValue {
return record.getBinary(idx);
}
+ @Override
+ public byte[] getVariantMetadata() {
+ return getVariant().metadata();
+ }
+
+ @Override
+ public byte[] getVariantValue() {
+ return getVariant().value();
+ }
+
+ private Variant getVariant() {
+ return record.getVariant(idx);
+ }
+
@Override
public void unpackArray(List<ColumnValue> values) {
InternalArray recordArray = record.getArray(idx);
diff --git
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTypeUtils.java
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTypeUtils.java
index 1b8046494d2..be457571f0b 100644
---
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTypeUtils.java
+++
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTypeUtils.java
@@ -43,6 +43,7 @@ import org.apache.paimon.types.TimestampType;
import org.apache.paimon.types.TinyIntType;
import org.apache.paimon.types.VarBinaryType;
import org.apache.paimon.types.VarCharType;
+import org.apache.paimon.types.VariantType;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -168,6 +169,11 @@ public class PaimonTypeUtils {
return paimonColumnType;
}
+ @Override
+ public PaimonColumnType visit(VariantType variantType) {
+ return new PaimonColumnType(Type.VARIANT);
+ }
+
@Override
public PaimonColumnType visit(ArrayType arrayType) {
PaimonColumnType paimonColumnType = new
PaimonColumnType(Type.ARRAY);
diff --git
a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonColumnValueTest.java
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonColumnValueTest.java
index 20062ff2f6e..c0651163c7d 100644
---
a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonColumnValueTest.java
+++
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonColumnValueTest.java
@@ -27,6 +27,7 @@ import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.serializer.InternalRowSerializer;
+import org.apache.paimon.data.variant.GenericVariant;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.BigIntType;
import org.apache.paimon.types.DataType;
@@ -36,6 +37,7 @@ import org.apache.paimon.types.MapType;
import org.apache.paimon.types.RowType;
import org.apache.paimon.types.TimestampType;
import org.apache.paimon.types.VarCharType;
+import org.apache.paimon.types.VariantType;
import org.junit.Assert;
import org.junit.Test;
@@ -92,6 +94,19 @@ public class PaimonColumnValueTest {
Assert.assertEquals(100L, values.get(2).getLong());
}
+ @Test
+ public void testGetVariantEncodedBytes() {
+ GenericVariant variant = GenericVariant.fromJson("{\"a\":[1,null]}");
+ PaimonColumnValue columnValue = new PaimonColumnValue(
+ GenericRow.of(variant), 0, ColumnType.parseType("v",
"variant"), new VariantType(), "UTC");
+
+ Assert.assertArrayEquals(variant.metadata(),
columnValue.getVariantMetadata());
+ Assert.assertArrayEquals(variant.value(),
columnValue.getVariantValue());
+ Assert.assertEquals(
+ ColumnType.Type.VARIANT,
+ PaimonTypeUtils.fromPaimonType("v", new
VariantType()).getType());
+ }
+
@Test
public void testReuseArrayElementsAcrossRows() {
ArrayType paimonArrayType = new ArrayType(new IntType());
diff --git
a/fe/fe-common/src/main/java/org/apache/doris/catalog/VariantType.java
b/fe/fe-common/src/main/java/org/apache/doris/catalog/VariantType.java
index c52d4a42304..df14e150c56 100644
--- a/fe/fe-common/src/main/java/org/apache/doris/catalog/VariantType.java
+++ b/fe/fe-common/src/main/java/org/apache/doris/catalog/VariantType.java
@@ -34,6 +34,7 @@ import java.util.stream.Collectors;
public class VariantType extends ScalarType {
private static final Logger LOG = LogManager.getLogger(VariantType.class);
+ public static final VariantType COMPUTE_V2_INSTANCE = new
VariantType(true);
@SerializedName(value = "fieldMap")
private final HashMap<String, VariantField> fieldMap = Maps.newHashMap();
@@ -83,6 +84,11 @@ public class VariantType extends ScalarType {
this.computeV2 = false;
}
+ private VariantType(boolean computeV2) {
+ this();
+ this.computeV2 = computeV2;
+ }
+
public VariantType(ArrayList<VariantField> fields) {
super(PrimitiveType.VARIANT);
Preconditions.checkNotNull(fields);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonUtil.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonUtil.java
index c0fda16413d..1a2d0511a4c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonUtil.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonUtil.java
@@ -26,6 +26,7 @@ import org.apache.doris.catalog.PartitionItem;
import org.apache.doris.catalog.PartitionKey;
import org.apache.doris.catalog.ScalarType;
import org.apache.doris.catalog.Type;
+import org.apache.doris.catalog.VariantType;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.UserException;
import org.apache.doris.common.util.TimeUtils;
@@ -342,6 +343,8 @@ public class PaimonUtil {
return ScalarType.createTimeStampTzType(tsScale);
}
return ScalarType.createDatetimeV2Type(tsScale);
+ case VARIANT:
+ return VariantType.COMPUTE_V2_INSTANCE;
case ARRAY:
ArrayType arrayType = (ArrayType) dataType;
Type innerType =
paimonPrimitiveTypeToDorisType(arrayType.getElementType(),
enableVarbinaryMapping,
@@ -360,8 +363,6 @@ public class PaimonUtil {
.map(field -> new
org.apache.doris.catalog.StructField(field.name(),
paimonTypeToDorisType(field.type(),
enableVarbinaryMapping, enableTimestampTzMapping)))
.collect(Collectors.toCollection(ArrayList::new)));
- case VARIANT:
- return Type.VARIANT;
case TIME_WITHOUT_TIME_ZONE:
return Type.UNSUPPORTED;
default:
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonUtilTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonUtilTest.java
index 236308d85a9..1e7f7f03b23 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonUtilTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/PaimonUtilTest.java
@@ -24,6 +24,7 @@ import org.apache.doris.catalog.PartitionItem;
import org.apache.doris.catalog.PrimitiveType;
import org.apache.doris.catalog.StructType;
import org.apache.doris.catalog.Type;
+import org.apache.doris.catalog.VariantType;
import org.apache.doris.datasource.NameMapping;
import org.apache.doris.datasource.metacache.paimon.PaimonPartitionInfoLoader;
import org.apache.doris.thrift.TPrimitiveType;
@@ -97,6 +98,15 @@ public class PaimonUtilTest {
Assert.assertEquals(14, type2.getLength());
}
+ @Test
+ public void testVariantMapsToComputeV2() {
+ Type type = PaimonUtil.paimonTypeToDorisType(
+ new org.apache.paimon.types.VariantType(), true, true);
+
+ Assert.assertTrue(type.isVariantType());
+ Assert.assertTrue(((VariantType) type).isComputeV2());
+ }
+
@Test
public void testTimestampWriteTypeMappingUsesDateTimeV2() {
RowType rowType = DataTypes.ROW(
diff --git
a/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out
b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out
new file mode 100644
index 00000000000..f8966577392
--- /dev/null
+++
b/regression-test/data/external_table_p0/paimon/test_paimon_catalog_variant.out
@@ -0,0 +1,26 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !desc --
+id bigint Yes true \N
+payload variant<PROPERTIES ("variant_max_subcolumns_count" =
"0","variant_enable_typed_paths_to_sparse" =
"false","variant_max_sparse_column_statistics_size" =
"10000","variant_sparse_hash_shard_count" = "1")> Yes true \N
+
+-- !full_variant --
+1
{"active":true,"age":18,"missing":null,"name":"alice","profile":{"city":"beijing","zip":100000},"score":98.5,"tags":["flink","paimon"]}
+2
{"active":false,"age":30,"extra":{"levels":[1,2,3]},"name":"bob","profile":{"city":"shanghai"},"tags":["doris"]}
+3 [1,"mixed",false,null,{"k":"v"}]
+
+-- !object_subpaths --
+1 alice 18 beijing true
+2 bob 30 shanghai false
+3 \N \N \N \N
+
+-- !null_and_missing --
+1 false true
+2 true true
+3 true true
+
+-- !root_array --
+3 1 mixed false null v
+
+-- !subpath_predicate --
+2 bob
+
diff --git
a/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
new file mode 100644
index 00000000000..e5a2cf79dec
--- /dev/null
+++
b/regression-test/suites/external_table_p0/paimon/test_paimon_catalog_variant.groovy
@@ -0,0 +1,91 @@
+// 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.
+
+suite("test_paimon_catalog_variant",
"p0,external,doris,external_docker,external_docker_doris") {
+ String enabled = context.config.otherConfigs.get("enablePaimonTest")
+ if (enabled != null && enabled.equalsIgnoreCase("true")) {
+ String minioPort =
context.config.otherConfigs.get("iceberg_minio_port")
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+ String catalogName = "test_paimon_variant"
+
+ sql """drop catalog if exists ${catalogName}"""
+ sql """create catalog if not exists ${catalogName} properties (
+ "type" = "paimon",
+ "paimon.catalog.type" = "filesystem",
+ "warehouse" = "s3://warehouse/wh",
+ "s3.endpoint" = "http://${externalEnvIp}:${minioPort}",
+ "s3.access_key" = "admin",
+ "s3.secret_key" = "password",
+ "s3.region" = "us-east-1",
+ "s3.path.style.access" = "true"
+ );"""
+ sql """use `${catalogName}`.`test_paimon_spark`"""
+ sql """set force_jni_scanner = true"""
+
+ explain {
+ sql "select * from variant_smoke order by id"
+ contains "paimonNativeReadSplits=0/1"
+ }
+
+ order_qt_desc """desc variant_smoke"""
+
+ order_qt_full_variant """
+ select id, payload
+ from variant_smoke
+ order by id
+ """
+
+ order_qt_object_subpaths """
+ select id,
+ cast(payload['name'] as string),
+ cast(payload['age'] as int),
+ cast(payload['profile']['city'] as string),
+ cast(payload['active'] as boolean)
+ from variant_smoke
+ order by id
+ """
+
+ order_qt_null_and_missing """
+ select id,
+ payload['missing'] is null,
+ payload['not_exist'] is null
+ from variant_smoke
+ order by id
+ """
+
+ order_qt_root_array """
+ select id,
+ cast(payload[1] as int),
+ cast(payload[2] as string),
+ cast(payload[3] as boolean),
+ cast(payload[4] as string),
+ cast(payload[5]['k'] as string)
+ from variant_smoke
+ where id = 3
+ order by id
+ """
+
+ order_qt_subpath_predicate """
+ select id, cast(payload['name'] as string)
+ from variant_smoke
+ where cast(payload['age'] as int) >= 20
+ order by id
+ """
+
+ sql """set force_jni_scanner = false"""
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]