Ravi Anne
02/24/2025, 7:56 PMResult::unwrap() on an Err value: NotYetImplemented("Casting from Map(Field { name: \"entries\", data_type: Struct([Field { name: \"key\", data_type: Utf8, is_nullable: false, metadata: {} }, Field { name: \"value\", data_type: Utf8, is_nullable: true, metadata: {} }]), is_nullable: false, metadata: {} }, false) to Map(Field { name: \"entries\", data_type: Struct([Field { name: \"key\", data_type: LargeUtf8, is_nullable: false, metadata: {} }, Field { name: \"value\", data_type: LargeUtf8, is_nullable: true, metadata: {} }]), is_nullable: false, metadata: {} }, false) not supported")`
sample test code is as follows
import datetime
import daft
import pyarrow as pa
from pyiceberg.catalog.sql import SqlCatalog
from pyiceberg.schema import Schema
from pyiceberg.types import (
ListType,
MapType,
NestedField,
StringType,
StructType,
TimestampType,
)
def local_catalog(tmpdir):
catalog = SqlCatalog(
"default",
**{
"uri": f"sqlite:///{tmpdir}/pyiceberg_catalog.db",
"warehouse": f"file://{tmpdir}",
},
)
catalog.create_namespace_if_not_exists("default")
return catalog
event_value_type = StructType(
NestedField(200, "name", StringType(), required=False),
NestedField(201, "timestamp", TimestampType(), required=False),
NestedField(202, "attributes", MapType(
key_id=203,
key_type=StringType(),
value_id=204,
value_type=StringType()
), required=False)
)
# Define the Iceberg schema for the trace data
events_iceberg_schema = Schema(
NestedField(field_id=1,name="events", field_type=ListType(2, event_value_type, required=True))
)
def events_table() -> tuple[pa.Table, Schema]:
events_array = [
[{"name":"e1","timestamp":datetime.datetime(2024, 2, 10),"attributes":{"a": "x1", "b": "y1"}}],
[{"name":"e2","timestamp":datetime.datetime(2024, 2, 11),"attributes":{"a": "x2", "b": "y2"}}],
[{"name":"e3","timestamp":datetime.datetime(2024, 2, 12),"attributes":{"a": "x3", "b": "y3"}}]
]
element_struct_type = pa.struct([
('name', pa.string()),
('timestamp', pa.timestamp('us')),
('attributes', pa.map_(pa.string(), pa.string()))
])
events_struct_type = pa.list_(element_struct_type)
events = pa.array(events_array, type=events_struct_type)
table = pa.table(
{
"events":events
}
)
return table, events_iceberg_schema
pa_table, schema = events_table()
catalog = local_catalog('/tmp/warehouse')
table = catalog.create_table_if_not_exists("default.test", schema)
df = daft.from_arrow(pa_table)
result = df.write_iceberg(table)
as_dict = result.to_pydict()
assert all(op == "ADD" for op in as_dict["operation"]), as_dict["operation"]
assert sum(as_dict["rows"]) == 3, as_dict["rows"]
read_back = daft.read_iceberg(table)
read_back.show()jay
02/24/2025, 7:57 PMRavi Anne
02/25/2025, 6:27 AMKevin Wang
02/25/2025, 6:30 AMRavi Anne
02/25/2025, 9:12 AMdaft.read_iceberg(table) fails with similar error. Would probably wait for the fix
df = daft.from_arrow(pa_table)
result = df.write_iceberg(table)Ravi Anne
03/03/2025, 3:46 PMdaft.read_iceberg(table) is working fine with daft 0.4.6Kevin Wang
03/03/2025, 5:32 PM