The Birth of Parquet
sympathetic.ink
sympathetic.ink
We never had scala thrift bindings, and the Java ones were awkward from Scala, so I wrote a thrift plugin in JRuby that used the Ruby thrift parser and ERb web templates to output some Scala code. Integrated with our build pipeline, worked great for the company.
I also wrote one era of twitter's service deploy system on a hacked up Capistrano.
These projects took a few days because they were dirty hacks, but I still got a below perf review for getting easily distracted, because I didn't yet know how to sell those company-wide projects.
Anyhow, about a month before that team kicked off Parquet, I showed them a columnar format I made for a hackweek based on Lucene's codec packages, and was using to power a mixpanel-alike analytics system.
I'm not sure whether they were inspired or terrified that my hack would reach production, but I like to think I had a small hand in getting Parquet kickstarted.
Partially explains how murder came to be https://github.com/lg/murder
> below perf review
That's some cheap bullshit. Fuck marketing-oriented corporate engineering.
For example, if the RHS of your cylinder has a slightly larger radius than the LHS the cylinder will commence turning to the left.
The upshot is the thick side picks up more snow than the thin side and the disparity in radii increases more rapidly still. The cylinder becomes a truncated cone which turns sideways and halts!
And in some cases the rollerballs get too tall for the bonding strength of the snow, so they break into parts that can restart the cycle if the slope is steep enough.
Best part of ClickHouse native data format is I can use the same ClickHouse queries and can run in local or remote server/cluster and let ClickHouse to decide the available resources in the most performant way.
ClickHouse has a native and the fastest integration with Parquet so i can:
- Query local/s3 parquet data from command line using clickhouse-local.
- Query large amount of local/s3 data programmatically by offloading it to clickhouse server/cluster which can do processing in distributed fashion.
- https://clickhouse.com/blog/apache-parquet-clickhouse-local-...
- https://clickhouse.com/blog/apache-parquet-clickhouse-local-...
I have a 15gb parquet file in a s3 bucket and I need to "unzip" and extract every row from the file to write into my database. The contents of the file are emails and I need to integrate them into our search function.
Is this possible to do without an unreasonable amount of RAM? Are there any affordable services that can help here?
Feel free to contact me (email in bio), happy to pay for a consult at the minimum.
https://duckdb.org/2024/03/29/external-aggregation.html
https://duckdb.org/2021/06/25/querying-parquet.html
If your DB is mysql or postgres, then you could read a stream from parquet, transform inline and write out to your DB
https://duckdb.org/2024/01/26/multi-database-support-in-duck...
And an unrelated, but interesting read about the parquet bomb
https://duckdb.org/2024/03/26/42-parquet-a-zip-bomb-for-the-...
- Write a pandas_udf function in pyspark.
- Parition your data into smaller bits so that the pandas_udf does not get too much data at the same time.
Something like:
```
from pyspark.sql import SparkSession
import pyspark.sql.functions as f
@f.pandas_udf(return_type=whatever)
def ingest(doc: pd.Series): # doc is a pandas series now
# your processing goes here -> write to DB e.t.c
pd_series_literal = Create a pd.Series that just contains the integer 0 to make spark happy
return pd_series_literal
spark = SparkSession.builder.getOrCreate()df = spark.read.parquet("s3 path")
df = df.repartition(1000). # bump up this number if you run into memory issues
df = df.withColumn("foo", ingest(f.col("doc_column"))
```
Now the trick is, you can limit how much data is given to your pandas_udf by repartitioning your data. The more the partitions, the smaller the pd.Series that your pandas_udf gets. There's also the `spark.sql.execution.arrow.maxRecordsPerBatch` config that you can set in spark to limit memory consumption.
^ Probably overkill to bring spark into the equation, but this is one way to do it.
You can use a normal udf (i.e `f.udf()`) instead of a pandas_udf, but apparently that's slower due to java <-> python serialization
I just wanted to mention that AWS Athena eats 15G parquet files for breakfast.
It is trivial to map the file into Athena.
But you can't connect it to anything else than file output. But it can help you to for example write it to smaller chunks. Or choose another output format such as csv (although arbitrary email content in a csv feels like you are set up for parsing errors).
The benefit is that there is virtually no setup cost. And processing cost for a 15G file will be just a few cents.
You create a new partitioned table/location from the originally mapped file using a CTAS like so:
CREATE TABLE new_table_name
WITH (
format = 'PARQUET',
parquet_compression = 'SNAPPY',
external_location = 's3://your-bucket/path/to/output/'
) AS
SELECT *
FROM original_table_name
PARTITIONED BY partition_column_name
You can probably create a hash and partition by the last character if you want 16 evenly sized partitions. Unless you already have a dimension to partition by.Our company (scratchdata.com, open source) is literally built to solve the problem of schlepping large amounts of data between sources and destinations, so I have worked on this problem a lot personally and happy to nerd out about what works.
- does this decompress to giant sizes? - can't you split the file easily, because it includes row-based segments? - why does it take months to solve this for one file?
The solutions I was able to put together using Dask and Spark and such were all insanely slow, they just got killed by Slurm without getting anywhere. In the end I went back to good ole' shell scripting with xxd to handle most of the heavy lifting. Finished in under an hour.
The appeal of these newfangled tools is that you can work with data sizes that are infeasible to people who only know Excel, yet you don't need to understand a single thing about how your data is actually stored.
If you can be bothered to read the file format specification, open up some files in a hex editor to understand the layout, and write low-level code to parse the data - then you can achieve several orders of magnitude higher performance.
But if you want to do some kind of grouping or for example pivoting rows to columns, I think you will still benefit from a distributed tool like Spark or Trino. That can do the map/reduce job for you in a distributed way.
15 GB is a real drag to do anything with. So it’s a real pain when someone says “I’ll just give you 1 TB worth of parquet in S3”, the equivalent of dropping a billion dollars on someone’s doorstep in $1 bills.
Thanks again for sharing your tool and insightful knowledge.
Depends a lot on what you want to do with the data of course, but if you want to filter and slice/dice it, my experience is that it is really fast and stable. And if you already have it on s3, the threshold for using it is extremely small.
Also: How does the tool you sell here solve the problem - the data is already there and can't be processed (15GB - funny that seems to be the scale of YC startups?)? How does a tool to transfer the data into a new database help here?
Maybe because the problem literally is "how to transfer this data into a database"
https://parquet.apache.org/docs/concepts/
Maybe the file in question only has one row group. Which would be weird, because the creator had to go out of their way to make it happen.
It might not be fast, but a quick 1-off solution that you let run for a while would probably do that job. There shouldn't be a need to load the whole file into memory.
```
def _get_duck_db_arrow_results(s3_key):
con = duckdb.connect(config={'threads': 1, 'memory_limit': '1GB'})
con.install_extension("aws")
con.install_extension("httpfs")
con.load_extension("aws")
con.load_extension("httpfs")
con.sql("CALL load_aws_credentials('hadrius-dev', set_region=true);")
con.sql("CREATE SECRET (TYPE S3,PROVIDER CREDENTIAL_CHAIN);")
results = con \
.execute(f"SELECT * FROM read_parquet('{s3_key}');") \
.fetch_record_batch(1024)
for index, result in enumerate(results):
print(index)
return results
```I ran the above on a 1.4gb parquet file and 15 min later, all of the results were printed at once. This suggests to me that the whole file was loaded loaded into memory at once.
To stream, fetch more batches.
What ddb does to get the batches depends on hand wavey magic around available ram, and also the structure of the parquet.
When I'm writing to postgres though I'm doing into entirely inside DuckDB with a `INSERT INTO ... SELECT ...` and that seems to stream it over.
``` df = pl.scan_parquet('tmp/'+DUMP_NAME+'_cleaned.parquet')
with open('tmp/'+DUMP_NAME+'_cleaned.jsonl', mode='w', newline='\n', encoding='utf8') as f: for row in df.collect(streaming=True).iter_rows(named=True): row = {k: v for k, v in row.items() if (v is not None and v != [] and v != '')} f.write(json.dumps(row, default=str) + '\n') ```
Have you tried downloading the file from s3 to /tmp, opening it with pandas, iterating through 1000 row chunks, pushing to DB? The default DF to SQL built into pandas doesn't batch the inserts so it will be about 10x slower than necessary, but speeding that up is a quick google->SO away.
Streams batches of rows
https://arrow.apache.org/docs/python/generated/pyarrow.parqu...
Edit — May need to do some extra work with s3fs too from what I recall with the default pandas s3 reading
Edit 2 — or check out pyarrow.fs.S3FileSystem :facepalm:
Of course then you might as well do all the processing you're interested in while the file is on your local disk, since it is probably much faster than the cloud service disk.
$200 to rent a machine that can run the naive solution for an entire day is peanuts compared to the dev time for a “better” solution. Running that machine for eight hours would only cost enough to purchase about a half day of junior engineer time.
Look at the parquet file metadata: use whatever tool you want for that. The Python parquet library is useful and supports s3.
How big are your row groups? If it’s one large row group then you will run into this issue.
What’s the number of rows in each row group?
I'm curious if anyone has experience with OpenLineage/Marquez or similar they'd like to share
https://bugs.debian.org/838338
These days Debian packaging has become a bit irrelevant, since you can just shove upstream releases into a container and go for it.
> apt-cache search hdf5 | wc -l
134
> apt-cache search netcdf | wc -l
70
> apt-cache search parquet | wc -l
0If you follow that link, you'll see polars and parquet are a large highly configurable collection of tools for format manipulations across many HPC formats. Debian maintainers possibly don't want to bundle the entirety, as it would be vast.
Might this help you, though?
https://cloudsmith.io/~opencpn/repos/polar-prod/packages/det...
Are there dark secrets?
it's very modern and perhaps hasn't been around long enough to have debian maintainers feel it's vetted.
for instance, documentation for Python bindings is more advanced than for Rust bindings, but the package itself uses Rust at the low level.
Python 3.10.12 (main, Nov 20 2023, 15:14:05) [GCC 11.4.0]
on linux
Type "help", "copyright", "credits" or "license" for more
information.
>>> import pandas
>>> pandas.read_parquet('sample3.parquet')
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File "/usr/lib/python3/dist-packages/pandas/io/parquet.py",
line 493, in read_parquet
impl = get_engine(engine)
File "/usr/lib/python3/dist-packages/pandas/io/parquet.py",
line 53, in get_engine
raise ImportError(
ImportError: Unable to find a usable engine; tried using:
'pyarrow', 'fastparquet'.
A suitable version of pyarrow or fastparquet is required for
parquet support.https://github.com/hangxie/parquet-tools/blob/main/USAGE.md#...
FWIW i think i share your general aversion to _not_ using packages, just for the tidiness of installs and removals, though i'm on fedora and macos.
https://pandas.pydata.org/docs/reference/api/pandas.read_par...
https://packages.debian.org/buster/python3-pandas
Perhaps more alarm is called for when this python+pandas and parquet does not work on Debian, but that is not the case today.
ps- data access in clouds often uses the S3:// endpoint . Contrast to a POSIX endpoint using _fread()_ or similar.. many parquet-aware clients prefer the cloudy, un-POSIX method to access data and that is another reason it is not a simple package in Debian today.
Polars and DuckDB are much better about memory management.
https://github.com/pandas-dev/pandas/blob/main/pandas/io/par...
says "fastparquet" engine must be available if no pyarrow
What crap. That's 'source-available', NOT open-source.
But at least co-option of terminology is an indicator of success.