Dexter Antonio
03/01/2025, 2:08 AMwrite_iceberg .Dexter Antonio
03/01/2025, 5:35 AMDexter Antonio
03/01/2025, 4:42 PMDexter Antonio
03/01/2025, 8:29 PMKevin Wang
03/01/2025, 8:36 PMKevin Wang
03/01/2025, 8:36 PMDexter Antonio
03/01/2025, 8:52 PMDexter Antonio
03/01/2025, 8:56 PMs3://<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 = 1000Dexter Antonio
03/01/2025, 8:56 PMDexter Antonio
03/01/2025, 9:16 PMKevin Wang
03/01/2025, 9:52 PMKevin Wang
03/01/2025, 9:55 PMDexter Antonio
03/01/2025, 10:07 PMdf = 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.Kevin Wang
03/01/2025, 10:17 PMColin Ho
03/02/2025, 2:00 AMimport 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"])Colin Ho
03/02/2025, 2:01 AMdaft.context.set_execution_config(_parquet_target_row_group_size_=16*1024*1024) Lowering it should decrease the amount of accumulating / concatenatingDexter Antonio
03/02/2025, 5:18 AMBasically 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
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?
Wouldn't raising, not reducing it, decrease the amount of accumulating / concatenating?Lowering it should decrease the amount of accumulating / concatenatingdaft.context.set_execution_config(_parquet_target_row_group_size_=16*1024*1024)
Colin Ho
03/02/2025, 7:40 PMWouldn'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.Dexter Antonio
03/06/2025, 3:57 AM--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.