I'm trying out `write_iceberg` .
# general
d
I'm trying out
write_iceberg
.
🔥 1
Still going,..
Is there anyway to speed this up?
The utilization of my EC2 instance is pretty low. Is there a way to speed this ingest up? Are there any parameters that I can tune? Should I try using a Ray cluster? I'm thinking about switching to EMR for this job
k
How large is your dataset?
Also curious about where your data is being stored. Is it on S3 or local disk or somewhere else?
d
Everything is stored in s3 as parquet files. The total dataset is around 4.2 Terabytes.
The dataset is hive partitioned, but I'm not utilizing that information when loading in the parquet files. This is the structure of the data in s3
Copy code
s3://<bucket_name>/<experiment_name>/<core_name>/<celltype_name>/<gene_name>.parquet
bucket count = 1 experiment_name count = 1 core_name count = 994 cell_type_name count = 3 gene_name count = 1000
Each parquet file has a little over 1000 columns
It takes 2 min and 13 s to write one core of data, so about 33 hours for all 994 cores. If there was a way to do the reads and writes in parallel by s3 prefix, that would really speed things up.
k
Hm yeah, 4.2 TB in 33 hours is only 0.28 gbps, which is not anywhere enough to saturate the network bandwidth of any ec2 machine. I'm not sure if we parallelize writes at the moment, @Colin Ho do you know?
Running the query on a cluster of machines using Ray will definitely improve the query performance quite a bit, but it seems that the single-machine performance can probably also be improved significantly for this query
d
Thanks for looking into it. Here is some more code. It is really straightforward.
Copy code
df = daft.read_parquet("ss3://<bucket_name>/<experiment_name>/**/*")
pyarrow_schema = df.schema().to_pyarrow_schema()
iceberg_schema = pyarrow_schema_to_iceberg_schema(df.schema().to_pyarrow_schema())  # my helper function
table = create_iceberg_table(iceberg_schema, "df_table",
                     partition_by="core_image_id")  # returns pyiceberg table 
df.write_iceberg(table)
I think part of what makes this ingest difficult is that each parquet is really wide, so maybe some soft assumptions are being violated leading to subpar optimization. The EC2 instance is being underutilized. It would be great if there is something that I can tune to better utilize the available resources. My sense is that a cluster wouldn't help if this single instance isn't being maxed out.
k
My guess is that our parquet write is the bottleneck. Our single-machine reads should be really good, even for weirdly shaped parquet files. If that is the case, multiple machines would increase write throughput. But anyway, looking to hear what Colin thinks
c
i think its because of the large number of columns and our reliance on pyarrow for writes. Basically what happens during a write is that we need to convert all 1000 columns into pyarrow arrays, and then call pyarrow for the write. This can cause contention on the GIL. Additionally, we accumulate data up till the row group threshold of 128 mb. This means we have to concatenate the accumulated data, which is memory intensive (especially with large number of columns). The write iceberg is parallel, and should utilize all cores, but is unable to because it needs to concatenate + convert a large number of columns to pyarrow. Here's an example script:
Copy code
import daft
import numpy as np

NUM_COLUMNS = 1000
NUM_ROWS = 1000000
NUM_UNIQUE_IDS = 10

# Generate and write data
data = {
    "id": np.random.randint(0, NUM_UNIQUE_IDS, size=NUM_ROWS),
}
for i in range(NUM_COLUMNS):
    data[f"col_{i}"] = np.random.randint(0, 1000000, size=NUM_ROWS)

df = daft.from_pydict(data).write_parquet("z")


# read back and write
df = daft.read_parquet("z2")
df.write_parquet("z", partition_cols=["id"])
I think the best way forward is for us to implement our own native writes. Though until then, one thing that can help is tuning the parquet target row group size:
daft.context.set_execution_config(_parquet_target_row_group_size_=16*1024*1024)
Lowering it should decrease the amount of accumulating / concatenating
d
Thank you for taking a look into this. Your explanation is really helpful. I double checked the size of the entire dataset and I have 2,949,500 parquet files with a total size of 7.4 TiB so the write speed is a little higher than originally estimated.
Basically what happens during a write is that we need to convert all 1000 columns into pyarrow arrays, and then call pyarrow for the write. This can cause contention on the GIL.
So the issue is that the conversion is not taking place in parallel right now? The code is something like
Copy code
pyarrow_dict = {}
for col in df.columns:
    pyarrow_dict[col] = pyarrow.from_series(df[col])
Can you point me to the place in the source code where this occurs?
Additionally, we accumulate data up till the row group threshold of 128 mb. This means we have to concatenate the accumulated data, which is memory intensive (especially with large number of columns).
If the row group threshold is fixed to 128 mb, why does this depend on the number of columns? Is it because the number of compactions is dependent on the row size?
daft.context.set_execution_config(_parquet_target_row_group_size_=16*1024*1024)
Lowering it should decrease the amount of accumulating / concatenating
Wouldn't raising, not reducing it, decrease the amount of accumulating / concatenating?
c
Sure, the code is here: https://github.com/Eventual-Inc/Daft/blob/main/daft/io/writer.py#L197 And no, conversion is not parallel.
Wouldn't raising, not reducing it, decrease the amount of accumulating / concatenating?
Daft reads parquet in a streaming fashion. The chunks that are streamed from parquet are generally much smaller than 128 mb. Therefore the writer will accumulate these chunks until
parquet_target_row_group_size
bytes of data, then concatenates the accumulated data and writes it out as a row group. E.g. if the threshold is 128MB, and the parquet reader emits chunks of 1 MB, then the writer will accumulate 128 chunks, concat them, and write.
d
Thank you for the explanation. I finally got this working with pyspark on EMR. The trick was just throwing a crazy amount of power at it. There was a lot of unneeded shuffling going on, because I couldn't get pyspark to infer my "hive-like" partitioning scheme.
Copy code
--conf spark.executor.memory=32g \
                --conf spark.driver.cores=16 \
                --conf spark.driver.memory=64g \
                --conf spark.executor.instances=65 \
                --conf spark.driver.maxResultSize=4G \
                --conf spark.task.maxFailures=5 \
                --conf spark.stage.maxConsecutiveAttempts=5 \
                --conf spark.sql.adaptive.enabled=true \
                --conf spark.sql.adaptive.coalescePartitions.enabled=true \
                --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkCatalog \
                --conf spark.sql.catalog.spark_catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog \
                --conf spark.sql.catalog.spark_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
                --conf spark.sql.catalog.spark_catalog.warehouse={warehouse_prefix} \
                --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
                --conf spark.sql.shuffle.partitions=1000 \
                --conf spark.default.parallelism=1000 \
                --conf spark.emr-serverless.driver.disk=200G \
                --conf spark.emr-serverless.executor.disk=200G",
Something that I couldn't get working with spark, or Ray would be to explicitly specify the partitioning in the data on disk. I'm not suggesting that you use query hints, but instead allow the specification of metadata about the structure of the data that can be used in query optimization.