hi, I am trying to create daft dataframe using fro...
# general
r
hi, I am trying to create daft dataframe using from_arrow. It is throwing below error. Daft version used is 0.4.3 `pyo3_runtime.PanicException: called
Result::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
Copy code
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()
j
yes I believe there was a separate but related issue reported for nested UTF8 to LargeUTF8 casting from last week as well... cc @Kevin Wang / @Jessie Young who might know more
r
cc: @Kevin Wang and @Jessie Young - is there any workaround ?
k
Hey! Sorry I was working on fixing this but had other things on my plate. A possible workaround would be to store string data in your Iceberg table as LargeString instead of String type
r
tried with large_string() type of pyarrow, and the below are successful, but
daft.read_iceberg(table)
fails with similar error. Would probably wait for the fix
Copy code
df = daft.from_arrow(pa_table)
result = df.write_iceberg(table)
daft.read_iceberg(table)
is working fine with daft 0.4.6
❤️ 1
k
Awesome! great to hear