Hi, I'm reading in new json files into my table, a...
# general
a
Hi, I'm reading in new json files into my table, and am reaching some pyarrow error:
Copy code
File "pyarrow/table.pxi", line 776, in pyarrow.lib.ChunkedArray.combine_chunks
  File "pyarrow/array.pxi", line 4777, in pyarrow.lib.concat_arrays
  File "pyarrow/error.pxi", line 155, in pyarrow.lib.pyarrow_internal_check_status
  File "pyarrow/error.pxi", line 92, in pyarrow.lib.check_status
pyarrow.lib.ArrowInvalid: offset overflow while concatenating arrays, consider casting input from `string` to `large_string` first.
Basically, I have an UDF which parses through the files and dumps the json content (atm simply) into the table. Could I either define "large_string()" as return datatype for the udf? Or is there even a better workaround ? Not very knowledgeable with pyarrow...
r
Offset overflows indicate that your concatenated JSONs are larger than 2gb, which causes an offset error. If you use PyArrow LargeString, I believe the offsets should be able to support strings larger than 2mb. Maybe try converting it to LargeString to give it a shot.
❤️ 1
j
^ yes, but 2GB actually 😆
r
Oh whoops sorry typo! Will edit.
j
@Ammar Bitar could you share some of the code you’re using with PyArrow? Our team should be able to provide some better suggestions with that code sample
Also you could consider using
batch_size
on the UDF to make each invocation of the UDF only operate on
N
number of rows.
a
Thanks guys! Yes, that's what I assumed the error to mean; I tried running with a small(er) batch_size, but I still encounter the error (and I made sure the max size of the json *batch_size (+ overhead) wouldn't exceed the 2gb limit... The code I'm using doesn't directly use PyArrow, but my issue basically goes as follows: 1. read delta table 2. read annotations a. download from minio object store (using boto3) b. dump (as json) into annotation column (for each row) c.
df = df.with_column("annotations", parser_udf(daft.col("path")))
Copy code
@daft.udf(return_dtype=daft.DataType.string(), batch_size=15)
class UDF_ANNParser:
    def __init__(self, endpoint, access, secret) -> None:
        self.s3_client = boto3.client(
            "s3",
            endpoint_url=endpoint,
            aws_access_key_id=access,
            aws_secret_access_key=secret,
            verify=False,
        )

    def __call__(self, paths) -> list:
        def fetch_annotations(s3_path: str) -> dict:
            try:
                if not s3_path.startswith("s3://"):
                    return {}

                bucket_key = s3_path.replace("s3://", "")
                if "/" not in bucket_key:
                    return {}

                bucket, key = bucket_key.split("/", 1)
                json_key = str(pathlib.Path(key).with_suffix(".json"))

                response = self.s3_client.get_object(Bucket=bucket, Key=json_key)

                data = response["Body"].read()

                annotations = json.loads(data.decode("utf-8").replace("'", '"'))
                return json.dumps(annotations)

            except Exception as e:
                return json.dumps({})

        results = []
        for p in paths.to_pylist():
            annotations_dict = fetch_annotations(p)
            results.append(annotations_dict)

        return results
1. write new deltatable
And here the explain output (not sure if we get any useful information from here?)
Copy code
== Unoptimized Logical Plan ==

* Project: col(path), col(size), col(num_rows), col(UUID), col(height), col(width), col(channels), col(sharpness), col(contrast), col(brightness), pyclass_udf(col(path)) as annotations
|
* PythonScanOperator: DeltaLakeScanOperator(None)
|   File schema = path#Utf8, size#Int64, num_rows#Int64, UUID#Utf8, height#Int16, width#Int16, channels#Int8, sharpness#Float32, contrast#Float32, brightness#Float32
|   Partitioning keys = []
|   Output schema = path#Utf8, size#Int64, num_rows#Int64, UUID#Utf8, height#Int16, width#Int16, channels#Int8, sharpness#Float32, contrast#Float32, brightness#Float32


== Optimized Logical Plan ==

* Project: col(path), col(size), col(num_rows), col(UUID), col(height), col(width), col(channels), col(sharpness), col(contrast), col(brightness), pyclass_udf(col(path)) as annotations
|
* PythonScanOperator: DeltaLakeScanOperator(None)
|   File schema = path#Utf8, size#Int64, num_rows#Int64, UUID#Utf8, height#Int16, width#Int16, channels#Int8, sharpness#Float32, contrast#Float32, brightness#Float32
|   Partitioning keys = []
|   Output schema = path#Utf8, size#Int64, num_rows#Int64, UUID#Utf8, height#Int16, width#Int16, channels#Int8, sharpness#Float32, contrast#Float32, brightness#Float32


== Physical Plan ==

* Project: col(path), col(size), col(num_rows), col(UUID), col(height), col(width), col(channels), col(sharpness), col(contrast), col(brightness), pyclass_udf(col(path)) as annotations
|   Clustering spec = { Num partitions = 2 }
|
* TabularScan:
|   Num Scan Tasks = 2
|   Estimated Scan Bytes = 665757536
|   Clustering spec = { Num partitions = 2 }
|   Schema: {path#Utf8, size#Int64, num_rows#Int64, UUID#Utf8, height#Int16, width#Int16, channels#Int8, sharpness#Float32, contrast#Float32, brightness#Float32}
|   Scan Tasks: [
|   {File {}}
|   ]
j
Are you then writing the data back out to deltalake? We do invoke PyArrow on the writing path...
a
Yes
But, I was hoping that batching on the UDF would take care of this down the road?
update - after updating to Dafts' latest version (0.4.2) this works 🤷 😄
j
Great! What version were you using before?
I'm guessing you were on something
< 0.4
, in the new version we actually completely changed our execution engine to be streaming based which will size the chunks we're writing out in a much better way
✅ 1
a
Yes, on the processing machine I still was using 0.3.11