:wave: given the current state of the Iceberg inte...
# general
g
👋 given the current state of the Iceberg integration, how much would work would it be use Nessie https://projectnessie.org/guides/ given its support for Iceberg REST?
k
I haven't played around with it yet but if it supports Iceberg REST, I believe it should work out of the box with PyIceberg, which is what we use for Iceberg integration!
g
yeah, I see this tutorial using pyiceberg https://www.dremio.com/blog/intro-to-pyiceberg/
so, I guess if we have nessie catalog setup, this might just work with no changes required?
k
Yep that would be my guess. If you try it out, do let us know how it goes!
g
cool will do!
🙌 1
@Kevin Wang small update, got our Nessie version upgraded to latest and just tested directly with
pyiceberg
and things seem to be working, will try daft now. only thing is that
pyiceberg
is not using S3 creds configured on Nessie server and expects client to provide them.
Copy code
catalog = load_catalog(
    "nessie",
    **{
        "uri": "<http://nessie:19120/iceberg/main>",
        "warehouse": "warehouse",
        "s3.endpoint": "<https://s3endpoing.com/>,
        "s3.access-key-id": os.environ["AWS_ACCESS_KEY_ID"],
        "s3.secret-access-key": os.environ["AWS_SECRET_ACCESS_KEY"],
    },
)

catalog.create_namespace("demo")
print(catalog.list_namespaces())

schema = Schema(
    NestedField(1, "id", IntegerType(), required=True),
    NestedField(2, "name", StringType(), required=False),
)

catalog.create_table("demo.test_pyiceberg_table", schema)

data = pa.Table.from_pydict(
    {
        "id": np.array([1, 2, 3], dtype="int32"),
        "name": ["Alice", "Bob", "Charlie"],
    },
    schema=schema.as_arrow(),
)

table = catalog.load_table("demo.test_pyiceberg_table")
table.append(data)

table.scan().to_pandas()
🙌 1
I tried writing on a branch and hit an error:
Copy code
daft.read_iceberg(catalog.load_table("demo.test_pyiceberg_table")).to_pandas()
   id     name
0   4     Dave
1   5     Earl
2   6    Frank
3   1    Alice
4   2      Bob
5   3  Charlie
6   1    Alice
7   2      Bob
8   3  Charlie

results = daft.from_pandas(
    pd.DataFrame(
        {
            "id": [7, 8, 9],
            "name": ["George", "Henry", "Isaac"],
        }
    )
).write_iceberg(catalog.load_table("demo.test_pyiceberg_table"), mode="append")

daft.read_iceberg(catalog.load_table("demo.test_pyiceberg_table")).to_pandas()

daft.exceptions.DaftCoreException: DaftError::External Internal IO Error when opening: data-services-experimentation/nessie-test/demo/test_pyiceberg_table_743042c2-8141-4e42-bcef-f3314215f728/data/95847077-617a-4afc-9469-a53891d0a432-0.parquet
results contains the single file name in the error above
data-services-experimentation/nessie-test/demo/test_pyiceberg_table_743042c2-8141-4e42-bcef-f3314215f728/data/95847077-617a-4afc-9469-a53891d0a432-0.parquet
the good news is that if I switch back to the
main
branch, the table is not impacted and can be queried
k
Hm. Does writing with PyIceberg work?
g
yeah, writing with pyiceberg works, something like below, and I can read back the updated results in daft
Copy code
table.append(
    pa.Table.from_pydict(
        {
            "id": np.array([7, 8, 9], dtype="int32"),
            "name": ["George", "Henry", "Isaac"],
        },
        schema=schema.as_arrow(),
    )
)
k
Gotcha. Is there anything more in the error message? I believe the next line should be
Details:
g
ah, it says
No such file or directory (os error 2)
k
Interesting, Daft seems to think that this is a local file, instead of on S3. Might be that we're not getting the path properly from S3. It's also a read error, wonder why it would happen on a write. Everything else other than the write works with Daft?
g
looking in
data
, the formatting is different for some reason? the pyiceberg file has the
00000-0-
prefix
Copy code
[2025-02-07 13:19:43 PST]   914B STANDARD 00000-0-2025b60b-6a19-42b2-b326-9d10e8d17d63.parquet
[2025-02-07 13:21:19 PST]   926B STANDARD 8a69900c-5ace-4f64-843f-91d2a4d73889-0.parquet
a more condensed example:
Copy code
catalog.create_table("demo.test_daft_table", schema)
table = catalog.load_table("demo.test_daft_table")

data = pa.Table.from_pydict(
    {
        "id": np.array([1, 2, 3], dtype="int32"),
        "name": ["Alice", "Bob", "Charlie"],
    },
    schema=table.schema().as_arrow(),
)

table.append(data)

data2 = pa.Table.from_pydict(
    {
        "id": np.array([4, 5, 6], dtype="int32"),
        "name": ["Dave", "Earl", "Frank"],
    },
    schema=schema.as_arrow(),
)

results = daft.from_arrow(data2).write_iceberg(
    catalog.load_table("demo.test_daft_table"), mode="append"
)
table.scan().to_pandas()
returns:
Copy code
id     name
0   1    Alice
1   2      Bob
2   3  Charlie
daft.read_iceberg(catalog.load_table("demo.test_daft_table")).to_pandas()
errors out
daft.exceptions.DaftCoreException: DaftError::External Internal IO Error when opening: data-services-experimentation/nessie-test/demo/test_daft_table_4ef4d783-bb1a-434b-8b6f-a464e7ce9f6b/data/8a69900c-5ace-4f64-843f-91d2a4d73889-0.parquet:
👍 1
the two files above are from this condensed example
so the the file is there, but something about parquet file formatting is making pyiceberg skip it when scanning, leading to the mismatch?
metadata and data (maybe helpful)
note:
table.scan().to_pandas()
failed, I think the table pointer it was looking at was stale
k
Was it working before and then it stopped?
g
it stopped working after I restarted the interpreter and tried to read the table again
k
Ah. Was the Daft error from before or after the restart?
g
before the restart, I am guessing with
table = catalog.load_table("demo.test_daft_table")
,
table
is maybe pointing at an older snapshot after the daft write, which is why the pyiceberg read looked like it worked.
adjusting example code, both reads fail for the same reason after daft write:
Copy code
data = pa.Table.from_pydict(
    {
        "id": np.array([1, 2, 3], dtype="int32"),
        "name": ["Alice", "Bob", "Charlie"],
    },
    schema=schema.as_arrow(),
)

catalog.load_table("demo.test_daft_table").append(data)
catalog.load_table("demo.test_daft_table").scan().to_pandas()

data2 = pa.Table.from_pydict(
    {
        "id": np.array([4, 5, 6], dtype="int32"),
        "name": ["Dave", "Earl", "Frank"],
    },
    schema=schema.as_arrow(),
)

results = daft.from_arrow(data2).write_iceberg(
    catalog.load_table("demo.test_daft_table"), mode="append"
)
catalog.load_table("demo.test_daft_table").scan().to_pandas()
daft.read_iceberg(catalog.load_table("demo.test_daft_table")).to_pandas()
k
Oh wait so the initial Daft read and write work, but the read after the write does not?
g
yeah,
pyiceberg
is looking for the file locally?
Copy code
>>> catalog.load_table("demo.test_daft_table").scan().to_pandas()
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/table/__init__.py", line 1471, in to_pandas
    return self.to_arrow().to_pandas(**kwargs)
           ^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/table/__init__.py", line 1438, in to_arrow
    ).to_table(self.plan_files())
      ^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/io/pyarrow.py", line 1456, in to_table
    if table_result := future.result():
                       ^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/monorepo/.css/bin/python/3.12.4+20240713/python/lib/python3.12/concurrent/futures/_base.py", line 449, in result
    return self.__get_result()
           ^^^^^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/monorepo/.css/bin/python/3.12.4+20240713/python/lib/python3.12/concurrent/futures/_base.py", line 401, in __get_result
    raise self._exception
  File "/Users/garrett.weaver/monorepo/.css/bin/python/3.12.4+20240713/python/lib/python3.12/concurrent/futures/thread.py", line 58, in run
    result = self.fn(*self.args, **self.kwargs)
             ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/io/pyarrow.py", line 1437, in _table_from_scan_task
    batches = list(self._record_batches_from_scan_tasks_and_deletes([task], deletes_per_file))
              ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/io/pyarrow.py", line 1517, in _record_batches_from_scan_tasks_and_deletes
    for batch in batches:
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyiceberg/io/pyarrow.py", line 1228, in _task_to_record_batches
    with fs.open_input_file(path) as fin:
         ^^^^^^^^^^^^^^^^^^^^^^^^
  File "pyarrow/_fs.pyx", line 789, in pyarrow._fs.FileSystem.open_input_file
  File "pyarrow/error.pxi", line 155, in pyarrow.lib.pyarrow_internal_check_status
  File "pyarrow/error.pxi", line 89, in pyarrow.lib.check_status
  File "pyarrow/_fs.pyx", line 1559, in pyarrow._fs._cb_open_input_file
  File "/Users/garrett.weaver/Library/Caches/pypoetry/virtualenvs/marketintel-enrichment-pipeline-0EJ9G2K5-py3.12/lib/python3.12/site-packages/pyarrow/fs.py", line 419, in open_input_file
    raise FileNotFoundError(path)
FileNotFoundError: /Users/garrett.weaver/monorepo/py/ds/marketintel/marketintel_enrichment_pipeline/data-services-experimentation/nessie-test/demo/test_daft_table_d0af3d92-cd4d-4b82-8277-02cbcfcd5a53/data/a1fa17c4-21bd-4aaf-8beb-de0f6b4303ba-0.parquet
k
That's interesting. So the PyIceberg table changes to point to the wrong thing after the Daft write
what version of Daft and PyIceberg are you using btw?
g
latest for both 0.4.3 / 0.8.1
👍 1
k
@Robert Howell would you be interested in investigating this?
g
oh I see it, the file written by daft is missing the
s3://
prefix in the manifest files:
Copy code
{'content': 0, 'file_path': 'data-services-experimentation/nessie-test/demo/test_daft_table_4ef4d783-bb1a-434b-8b6f-a464e7ce9f6b/data/8a69900c-5ace-4f64-843f-91d2a4d73889-0.parquet', 'file_format': 'PARQUET', 'partition': {}, 'record_count': 3, 'file_size_in_bytes': 926, 'column_sizes': [{'key': 1, 'value': 93}, {'key': 2, 'value': 101}], 'value_counts': [{'key': 1, 'value': 3}, {'key': 2, 'value': 3}], 'null_value_counts': [{'key': 1, 'value': 0}, {'key': 2, 'value': 0}], 'nan_value_counts': [], 'lower_bounds': [{'key': 1, 'value': b'\x04\x00\x00\x00'}, {'key': 2, 'value': b'Dave'}], 'upper_bounds': [{'key': 1, 'value': b'\x06\x00\x00\x00'}, {'key': 2, 'value': b'Frank'}], 'key_metadata': None, 'split_offsets': [4], 'equality_ids': None, 'sort_order_id': None}
It would interesting to test this on just iceberg rest api to determine if it is that or specific to nessie
k
Yeah that's true. It might be an issue with Daft Iceberg writes in general? We should definitely look into this
g
I have been doing a lot of S3 Iceberg writes with daft and only encountered this while testing rest api, so hopefully just something around that 🤞
r
Ok, could either of you two please file an issue with these details so we don't lose the thread? Thank you.
g
r
Awesome, thank you!
g
stepping through, the resulting data files returned after collect have no prefix:
Copy code
data_files[0].file_path
'data-services-experimentation/nessie-test/demo/test_daft_table_d0af3d92-cd4d-4b82-8277-02cbcfcd5a53/data/e94293a7-6c1c-4c9b-b08d-a3b370963d88-0.parquet'
left a note about this one https://github.com/Eventual-Inc/Daft/issues/3783#issuecomment-2651602399. I am pretty motivated to get this one working and can try to help, but depends on where this change would be made with respect to how long it would take me to try to help 😂