Ammar Bitar
01/09/2025, 5:39 PMFile "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...Raunak Bhagat
01/09/2025, 5:46 PMjay
01/09/2025, 5:49 PMRaunak Bhagat
01/09/2025, 5:49 PMjay
01/09/2025, 5:50 PMjay
01/09/2025, 5:53 PMbatch_size on the UDF to make each invocation of the UDF only operate on N number of rows.Ammar Bitar
01/10/2025, 9:51 PMdf = df.with_column("annotations", parser_udf(daft.col("path")))
@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 deltatableAmmar Bitar
01/10/2025, 10:02 PM== 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 {}}
| ]jay
01/10/2025, 10:04 PMAmmar Bitar
01/10/2025, 10:11 PMAmmar Bitar
01/10/2025, 10:12 PMAmmar Bitar
01/16/2025, 7:52 AMjay
01/16/2025, 8:03 PMjay
01/16/2025, 8:04 PM< 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 wayAmmar Bitar
01/17/2025, 9:33 AM