← All transcripts

S3 to GPU at 60Gbps: Why Parquet is Bottlenecking Your Data Transcript, AI Summary & Key Points

InfoQ · 3 days ago · Science & Technology · 49:38 · EN

Watch on YouTube

Answer

Parquet pipelines bottleneck GPU workloads through NVMe staging, CPU decompression, and host-to-device memory movement. Vortex reduces those costs by pruning data before transfer, computing on compressed arrays, and streaming data toward the GPU at up to 60 Gbits per second.

AI Summary

Vortex is an open-source columnar file format maintained under the Linux Foundation that targets GPU analytics and training workloads. It streams compressed data from S3 toward GPUs, avoids mandatory NVMe staging and CPU decompression, prunes segments with layouts and zone maps, and supports compute on compressed arrays through cascading encodings. Its scan layer coalesces byte ranges, coordinates I/O with CPU or CUDA computation, and can use pinned buffer pools, KTLS, custom HTTP clients, or RDMA to reduce memory copies. The demonstrated pipeline sustains 60 Gbits per second end to end, while Vortex integrates with Arrow-compatible data, DataFusion, DuckDB, and Polars.

Key Points

  • A Vortex file can represent 4K 60fps RGB frames as three columns, with around 8 million pixels per frame, streaming roughly 13 Gbits per second from S3 through the CPU to the GPU.
  • GPU workloads pay a movement tax when data moves from S3 to NVMe, into RAM, through CPU decompression, and across PCIe to the GPU.
  • GPU workloads also pay a decision tax when changing a training data mix or filter requires reprocessing, retokenizing, reshuffling, and reloading the data.
  • Vortex separates logical types from physical types, uses an extensible minimal file specification, and performs work as late as possible to reduce materialization and memory copies.
  • Vortex uses lightweight cascading encodings rather than block compression that loses the meaning of the underlying data, allowing random access, aggregation, and other computation on compressed arrays.
  • Vortex S3-to-GPU scans are around 30 times faster than Parquet in the cited comparison, while random access is around 100 times or more faster.
  • The file structure contains a fixed-size postscript, a footer listing arrays and layouts, and aligned segments containing most of the file data.
  • Segment alignment lets receiving buffers be allocated for SIMD operations or GPU kernels without an additional copy before computation.

Tools & resources

8 items

ANo. 2144
AIAINotes.us Tool

Apache Arrow

Open source · apache/arrow

Apache Arrow is an open-source, language-independent universal columnar in-memory format and multi-language toolbox for fast data interchange and in-memory analytics. It provides a standard columnar representation for flat and nested data, an IPC format, Flight RPC, zero-copy reads, and libraries and development building blocks across languages including C, C++, .NET, Go, Java, JavaScript, Julia, MATLAB, Python, R, Ruby, Rust, and Swift. The project is an Apache Software Foundation effort developed by a distributed developer community, with readers and writers for formats such as Parquet and CSV.

Mentioned in
2 videos
Kind
Other
ANo. 4585
AIAINotes.us Tool

Apache DataFusion

datafusion.apache.org

Apache DataFusion is an extensible query engine written in Rust that uses Apache Arrow as its in-memory format. It provides SQL and DataFrame APIs, a query planner, and a columnar, streaming, multithreaded, vectorized execution engine for partitioned data sources. It supports CSV, Parquet, JSON, and Avro out of the box, and can be extended with custom data sources, query languages, functions, operators, and other workload-specific components. DataFusion is developed as an Apache Software Foundation project and is available as libraries and binaries for building database and analytics systems.

Mentioned in
1 video
Kind
Other
CNo. 0239
AIAINotes.us Tool

CUDA (Compute Unified Device Architecture)

nvidia.com

CUDA is a parallel computing platform and programming model developed by NVIDIA that enables general-purpose computing on NVIDIA GPUs. It includes the CUDA Toolkit (compilers, libraries, and developer tools) and APIs for C, C++ and Fortran to accelerate compute-intensive applications in areas such as scientific computing, HPC, graphics, and machine learning.

Mentioned in
3 videos
Kind
Other
DNo. 4586
AIAINotes.us Tool

DuckDB

duckdb.org

DuckDB is an in-process SQL OLAP database management system for analytical workloads. It uses a PostgreSQL-inspired SQL dialect, can run from edge devices to servers, and reads and writes formats including CSV, JSON, Parquet, Iceberg, and Vortex both locally and in object storage. Extensions add functions and format support, while native clients provide APIs for Python, Go, Java, Node.js, C, C++, R, Rust, and ODBC. DuckDB is open source and MIT-licensed, with governance by the independent DuckDB Foundation.

Mentioned in
1 video
Kind
Other
KNo. 4588
AIAINotes.us Tool

KTLS

In the AINotes directory

KTLS is a Linux kernel facility for decrypting TLS network data in the receive buffer, eliminating a memory copy into user space.

Mentioned in
1 video
Kind
Other
PNo. 4587
AIAINotes.us Tool

Polars

polars.dev

Polars is a Rust-backed data-processing engine that uses Apache Arrow's columnar in-memory format, multithreaded execution, SIMD vectorization, and a lazy query optimizer. Its Python API supports DataFrames, declarative expressions, joins, aggregations, window functions, schema-aware ingestion, and Parquet scanning. The lazy API applies predicate and projection pushdown before execution, while eager operations load data immediately; the page also describes out-of-core processing and PySpark migration workflows. The video identifies Polars as a data-processing engine that can interoperate with Vortex and other columnar-data systems.

Mentioned in
1 video
Kind
Other
RNo. 0706
AIAINotes.us Tool

Rust

Open source · rust-lang/rust

Rust is a systems programming language and toolchain developed by the Rust project. Its compiler, standard library, documentation, package manager and build tool Cargo, formatter rustfmt, linter Clippy, and editor support are maintained in the main repository. The language uses a rich type system and ownership model to enforce memory and thread safety at compile time, while targeting fast, memory-efficient software for critical services, embedded devices, and integration with other languages. Rust is primarily distributed under the MIT and Apache 2.0 licenses, with portions under BSD-like licenses.

Mentioned in
11 videos
Kind
Other
VNo. 4582
AIAINotes.us Tool

Vortex

In the AINotes directory

Vortex is an open-source columnar file format for streaming compressed analytical and machine-learning training data from object storage such as S3 to CPUs and GPUs. It separates logical and physical types and uses layouts, zone maps, cascading encodings, and byte coalescing to reduce unnecessary I/O, decompression, and memory copies; its encodings include dictionary arrays, run-end encoding, and bit-packing, allowing computation on compressed data. The format is designed for zero-copy execution with Apache Arrow and integration with DataFusion, DuckDB, and Polars, with scan paths that can use buffer pools, kTLS, and RDMA to move data toward CPU or GPU execution.

Mentioned in
1 video
Kind
Other
🔒 Full analysis locked

Unlock more videos and the full analysis

A credit unlocks one video's full analysis for good — the build steps, the tools and how each was used, the methods behind every use case. Pro opens the whole library instead, and raises how many videos you can analyse a day.

Unlock full analysis — free

Transcript

Searchable transcript of S3 to GPU at 60Gbps: Why Parquet is Bottlenecking Your Data — InfoQ (49:38). Search for a phrase, then click its timestamp to jump straight to that moment in the video.

Captions sourced from the original video on YouTube, published by InfoQ. The video, its captions and all related intellectual property remain the property of their respective owners; AINotes claims no ownership. Provided for research, accessibility and search — see the Transcript Notice and Copyright Policy.

00:12 Hi and welcome. Um, but so before I get into anything, I want to show you this video. Well, not specifically like what's playing in the video, but how we produce this video because this is not backed by a regular MP4 file in some disk. But this is actually a visualization of a file scan. What I mean by this is we have a vortex file. Vortex being a column file format.

00:37 It has three columns, one per color. And each each column would then have a row per frame of this video. And each frame would include around 8 million pixels because it's a 4K video. And then we then process this data which is three columns around 8 million pixels per frame, 60 frames per second. So around 13 Gbits per second from S3 through the network card over the CPU and to the GPU.

01:04 And the GPU is seeing this data this fast. But we want to see what's happening in the GPU. We want to visualize that. But we can't down link that much data to our laptops. So because it's already in the GPU, I'm using the GPU itself to re-encode it into H264, wrap it into TCP buffers, send it over my SH tunnel, and you can display and see what's happening.

01:23 And now have a look at this. Well, this is the same video, but this doesn't have any post-processing on it. This is actually a visualization of a column project column pruning or projection pruning where we only select in our projection expression one column and only the bytes that are associated with that column are flowing through the same pipeline keeping the same bandwidth uh and then ending up with the video that you were seeing before.

01:50 So welcome to my talk. Today I want to talk about how this is possible and the technologies that enable us to do this. I'm Our I work at Spiral DB. I help maintain Vortex which is an open source file format um that is under Linux Foundation and I want to talk about what makes Vortex special. Well, I'll start with why we need this speed in the first place.

02:15 I'll then move on to Vortex itself, explain some of the um main topics around vortex layouts and arrays and the scan itself um and then end up with um what is possible. So to start with the problem um why we need the speeds well because you have some GPUs that you hard and you want to use them. GPUs are expensive and then each time they're idling it's expensive and it's opportunity cost of not doing valuable work.

02:44 So I want to tell this like we have two types of taxes that we pay when utilizing GPUs. The first one is a movement tax. So getting the data from some source that you have into the GPU is not trivial and it's not really cheap. If you use web data set for training or for example mosaic ML's um streaming data set then they go more or less these steps they you fetch for example let's say your data is an S3 you fetch from S3 into the disk into NVME from NVME you read compressed bytes into your RAM you decompress on the CPU

03:15 and then send those decompressed bytes over the hosted device link PCIe to the GPU this not only has many steps that you have to follow True. If you zoom in into the uh supported bandwidth per connection, you can see that we are bottlenecked mainly by CPU decompression and the NVME throughput. And the reason why they all persist into NVME is the main consumer of these pipeline is PyTorch and PyTorch is Python.

03:43 Python doesn't like multi-threading it. So therefore, we use multiprocessing. Each worker doesn't want to step into each other's so they want to dduplicate the download. So they persist the the downloads on the disk. the shards on the disk. But ideally we want this pipeline to look more like that where you fetch from S3 eliminating the bandwidth bottlenecks through the network card speeds and the PCIe speeds to the GPU and then further you want to use the right compression that enables you to decompress in the GPU.

04:17 the parallel compressions and GPU can decompress much faster than CPUs and you're not only eliminating the bottlenecks you had in NVME and you're streaming much much faster now you're streaming maybe 50 Gbits per second but those bytes are compressed so if you pick the right compression you're effectively this pipeline runs through 10 times faster than the actual compressed speed that was movement and there's another tax I want to talk about which is decision it it This is when you want to iterate uh on your training

04:49 workload. You you've loaded from your your your source into the GPU. Then maybe you did some pre-processing, you are training, then you want to evaluate the results. And what you want to see then try for example is a new data mix or you want to change the curriculum. You want some sort of new filter on this data. Well, with the with the current uh tools we have, you have to reprocess the entire data, re probably retokenize, reshuffle to get the right order and then load again to your source, which might be S3, might be

05:20 somewhere else. But this is slow. And if you look at the hypothetical pie chart of your percentage of time spent, you want to maximize the gray where you're thinking and the green where your GPU is working and everything else is overhead. And you want to turn this into that where you load from your source into the GPU without wasting much time. And then once you evaluate and decide you want to change something in your input source, then you just change your scan.

05:47 You just give a different query, different filter, you don't need to reprocess. And then you immediately can start training for your next iteration. So how does vortex uh using a different file format enable all this? Well, vortex to give a quick overview of what it is. It's a columnal file format. It is uh similar to parket, but it's also much different than paret.

06:08 Um, it has a an extensible spec, a minimal spec that you can plug in into various points. It has a separation of logical types versus physical types. That's like the one of the main differences from parquet is that vortex doesn't have a notion of tying the types themselves like integers into how they're encoded on disk. It's up to you. is decoupled.

06:29 It does as little as possible as late as possible saving a lot of mem copies for materializing intermediate results and like leveraging from query optimizations and then it can push through most of the projections, filters and aggregation from your compute engine into the file format itself. Also, the way it does compression is different. Um, it doesn't use block compression that loses the meaning of the underlying data.

06:56 It does use lightweight uh and cascading encodings um that can can also give a similar compression ratio but more importantly they allow you to do compute on the compressed data. You can do random access you can do aggregations and some of these compute can be even faster than doing it on the decompressed data itself. Uh to give some concrete numbers uh versus parquet this S3 to GPU scans that we were showing is around 30 times faster.

07:24 Um and with thanks to the uh encodings and the the compression codeex vortex has it's around 100 times or more faster than parquet in random access. So what does the file look like? Well um we will approach this from a file reader perspective. There's a post script that's fixed size. You read it. It has a lot of pointers into other sections of the file.

07:45 Um there's a footer. The footer itself is a legend. It's a bill of materials of what you expect to see in this file. You have a list of arrays and you have a list of layouts. Arrays are how you encode the data itself in vortex. And layouts are like the logical plan. They are the IO layer of of pruning. And we will dive into that later. But these are showing you the main plug-in point.

08:09 So if you have a new array of a new layout, you just add it here and you have the logic in your reader, then vortex will delegate to your reader to do what what you want to do. And thirdly, and maybe most importantly, it's the list of segments that you expect to see in this file. And segments are the majority of the file. They're array data themselves.

08:27 They come with triplets. So offsets, length, and alignment to their pointers, but they also have alignment. Alignment is important to bake in the file. We think because if you have the alignment you know that if you want to use SIMD operations or GPU kernels when you load this well if you know the alignment beforehand in your receiving buffer you you allocate that buffer with the right alignment so you don't need to copy just to dispatch the compute on top and I've talked about arrays and layouts but what are layouts

08:59 layouts enable you to prune those segments as quickly as possible given a filter and projection they are like logical plan they are also very minimal. Well, this is the entire flat buffer definition. The only two maybe important points I want to touch is the last two fields. It has a children which is an array of itself. So, it's a tree and each node in the tree has segments.

09:21 Segments is a list of integers and those integers point to the offsets of segments and we know from the footer where they are in terms of pointer uh like exact bite ranges. So layout is able to tell each node in layout tells I am responsible of reading these segments and give me a filter and I'll tell you which of these segments I'm responsible of you have to read.

09:44 And to give an example on how layouts work. Let's start with a with a schema for example. So we have the root layout. It has two columns a name string and age integer. So the strruct layout only tells that I have a child per column. So you expect it to have two different children one per column and each of these is a zoned layout itself. Zoned layout is interesting because it not only contains the data as segments it also contains another segment that is a summary statistics of of this of the of the data that you

10:16 expect to read. That's also a segment itself. But when you read that segment it's a zone map. A zone map is al is a vortex array by itself. But it does have a row per row range in the actual data and has columns of uh statistics like min, max etc. And the data is then just a concatenation of multiple segments. So it's a way of of breaking it down into smaller arrays when you want to read them.

10:40 How how does this work then? Well, like say you have a query, you want to get the name column, but you only want to get the rows where the age column is greater than 30. Well, you start from the root. root knows it has a projection and a filter in each column. So it delegates into each each zone layout readers. The age will receive a filter and it will then initially load the zone map.

11:02 Check the zone map, check the filter, find which segments in this data passes the filter. Let's say the last segment passes the filter. So then it returns the last last segment. We also get the name column to return the last segment because that's the actual data we need. Do the actual filter on the actual rows and return the rows that we need. So by reading the zone map and only 30% of the segments you can you can satisfy this query.

11:27 Maybe this is a more detailed view. So the sensor is the zone map. That was a segment. We read it. We des serialized it. It's an array. It has one row per row range. And we have the filter. So by just looking at the max value, you can prove that only the last row here can have any rows that match this filter. All the others you don't need. You read the entire chunk.

11:48 That's the that's the vortex atomic unit of reading is a segment inside. And then what's neat about vortex is the compression we have is slicable. So you don't need to decompress everything. You slice it first and then you only decompress the rows that match the row the row range we have it. So last 8k rows that is decompressed. So that is like doing as little as possible and as late as possible.

12:11 And the difference between chunk granularity and row granularity is also I think important to mention because these serve different purposes. Chunk granularity is an IO concern. It deals with what type of IO source you have, what its bandwidth and latency characteristics. So you can tune them in vortex and the zone uh granularity itself is a compute concern because you want to depending on your architecture you might have different wider CID registers that you want to use and those play well with different uh lengths

12:42 of data that you have. So you can tune them as to your liking uh with vortex. So we talked about reading the segments and layout is helping us prune to the segments that we want. But what do each segment contain? They contain a serialized most likely compressed vortex array. And a vortex array I want to also explain uh with an example. Let's say we have this array.

13:05 It's an array of arrays. So it's the element is a list itself. And the elements inside the inner lists are strings. So we want to encode them and compress them as as as most as possible while keeping the meaning. So we can start with a list array that's similar to arrow. um it basically flattens the values and saves the flattened values as a values child and then it keeps the inner brace locations because it removes them while flattening it loses the information.

13:32 So it keeps that information into the offsets child. So any two consecutive offsets would point to the inner braces of an inner list location. Um and you can see this is also a tree right now we have we have the list array but it's two children are also arrays themselves. So we can we can go on further we can compress further. For example, the values child has only three unique elements.

13:52 So we dictionary encode them. Dictionary encoding gets the unique values and stores them as a children. So like a values child has one city appearing once only. But you then encode the repetition of the of the original array as the indices that point to the values child. So then you decoupled the values themselves and the repetition into separate children right and you can continue further because the codes now have continuous repeated values which are called runs.

14:22 So you can run length encode them or in vortex more specifically run end encode them meaning you get the values that repeat save them as a child zero one and two are the values of each run here and then you also encode which offsets in the original array these runs end. So the actual indices that end. So with this we have we have the the the string values as a as a child but we also have multiple integer values that we have and we can see that all of these values are really small.

14:54 The values child of run end they're all smaller than four. Four is a power of two second power of two. So you can encode each of them in two bits. The ends and the offsets are all less than eight. So you can use three bits to encode them. So you bit pack all of this integer data you have and end up with the actual bits on disk like that. So that means this is total 30 bits or so.

15:17 So this is less than one one integer. It's less than four bytes. And then you using this information you are storing all the repetition and nested structure of your input array and you only store the unique elements once. We can further compress this. We have fssts for example is an algorithm that extracts the sub common substrings and then further divides into into arrays.

15:41 uh but for this example I think this is this gets the point and I want to I want to kind of separate these two things to get in in two two ways. Well, you have the arrays which is the which is a structure which is how you encode this segment. Uh and then it is it is important because this is what you would use to dispatch compute and on the right hand side you have the data.

16:06 I want to think of these two as the control plane and the data plane and that will be important in the scan which we'll go into later but before I want to mention one more thing ha having this array tree on hand and having these data uh as pointers somewhere in some location you can dispatch compute without decompressing um for example I want to get a random read I want to get the second element of the first list well if you get that query into the list array list array knows that it can binary search the offsets child

16:37 to find the actual location then does propagate that to the flattened values child because it now knows the actual location. Dictate can do the same for codes and then substitute with the right value before returning. Um and also I mentioned some compute you can run on compressed data might end up being faster than running the data running the compute on decompressed arrays like summing for example.

17:00 Summing requires on a decompressed array iterating through each element and summing the values. Well, if you have a run end array, that's just multiplying with the lengths of each run and then per run you are summing the total uh like result of the multiplication. So, it's a lot faster um to to run the compute on compressed data. Okay, back to the separation.

17:22 So, the control plane shows you the amount of information you need in a compute engine so you can dispatch work on the data itself. And the data plane is the information you have to transfer to your compute engine. So these um serve different purposes and we don't want them to step in each other's toes. How do we do this? Uh well we use a vortex has a scan that tries to maximize this um throughout.

17:50 Um the way it works is it tries to maximize the IO engine's work and the compute engine's work to be busy all the time. And then scan is the orchestrating layer. So the orchestrating itself also wastes compute and has to um dispatch work decompress well reconstruct the array trees and then follow the array tree logic all the way to the end. Um the IO layer then gets some bytes to the compute source returns the pointers.

18:20 The scan then dispatch work on those pointers and well a small detail but compute also can request more segments because for example in the query that we initially did uh select name where age is greater than 30 you don't know if you need the name column before executing the age filter. So it could just be pruned alto together so you don't end up getting those segments.

18:43 So say we have a query in the scan and then we have the layout in the scan so we know which segments we need to fetch. Those are translated into bite ranges and these are bite ranges that you want to get from your source. So um each individual byte range might have a latency and depending on what source you're reading you don't want to dispatch each individual by ranges like by itself.

19:10 So we have a coallesing layer on top and coalesing means if the two bite ranges are close enough sufficiently you merge them and then have a single bite range. You get extra bytes in between but because you pay the latency once and get the entire throughput it kind of uh works on your on your benefit. But this is uh source dependent. NVME is happy with some calls in configuration but S3 for example might require a lot more.

19:34 uh SG can you can you can do 16 megabyte segments um because it's huge through triput but huge latency. So anyway once you get the segment range uh requests and then have the bite ranges you then dispatch work on them. Well, if it's a CPU scan, this work is done when you have the bytes on RAM. And then the work you dispatch is regular Rust functions like slice, projection, filter, any compute that the the input query wants you to do and that was CPU.

20:04 But um most of the things stay the same when you are scanning to the GPU. The scan orchestration stays the same. The coallesing layer stays the same because we still need those B ranges. But now the IO plane uh has to ship these buffers all the way to the GPU and return pointers to the GPU. That's when its job is done. And the compute now instead of dispatching Rust functions, it does dispatch CUDA kernels.

20:31 But everything else is the same. We have a buffer pool in the middle. And I want to explain why we have a buffer pool in the middle next. So this is the lifetime of a bite range um through a scan from S3 to the GPU. And initially when you request bytes from S3, it lands on the kernel receive buffer. This is a direct uh memory copy. So it's not an actual copy.

20:57 It's the first place that these bytes land from the network card is this kernel receive buffer. Then it immediately copies the cipher text still in TLS encrypted form into the userland one to one. Then we are in now or in our HTTP libraries land where we use ROSTLS and request RSTLS decrypts it into a plain text copies. That's one copy and then it's still TCP.

21:20 We want to get the HTTP body out of it. So the request layer D frames it gets the body into its own buffer. Then finally it arrives to our own vortex land. And because these yellow buffers are all streaming because it's a multiart get request, we have to aggregate them because we want the segment. We want a continuous bite range. So we aggregate them.

21:40 So we have to copy one more time. But that's the only copy we do. We receive these chunks as a streaming fashion and then we accumulate them until the request is done. And immediately after this is done, we dispatch the CUDA. We do CUDA async copy and expect it to land on the GPU. But unbeknownst to us, internally, CUDA does another copy. It does another copy because it requires some alignment uh rules um from from us.

22:04 Basically um the buffers you can copy to the GPU from the host has to be uh pinned and page locked meaning it should be outside the operating systems virtual page system. So then the copy if if the buffer gets paged out during the copy uh GPU uh doesn't seg and stays happy. So one obvious way to well before before before diving into how we can reduce the copies I want to mention one more thing the mem copy cost uh is what's important to us it's the unit is a core per gigabit of bandwidth you can sustain and it's

22:43 important because this core uh budget is also used by the orchestrating layer in scan and also compute if you're doing the compute on the CPU so it has the opportunity cost um if because you have limited CPU if you're just wasting it on me copies um all around then you slow down your own bandwidth. It's like a double-edged damage. So you want to reduce this as much as possible.

23:05 So how can we reduce this? Well, initially um we start with the CUDA bounce buffer because we know the requirements. We can um we can satisfy those requirements on our on our own buffer that we accumulate the HTTP buffers into. Um so we have a pinned buffer that we use and then we then call CUDA async copy on it. it doesn't bounce use a bounce bounce off and go goes directly into the GPU and further because allocating these are expensive we keep a buffer pool around to offset the amortize the the allocation cost and

23:38 the way we reuse these buffers is when you dispatch the CUDA copy you append the CUDA CUDA event after the after that uh CUDA stream append an event and once that event fires you know everything that comes before it is complete so then you can safely claim the buffer because that means the transfer is complete and you can reuse that same buffer for for next uh copies in your pipeline.

23:58 So this reduces the cost a bit but we need to then look to the left to the HTTP land to to get rid of some of the copies. One thing you can do here is use something called if you have a sufficiently high Linux version you can use a KTLS kernel site TLS decryption. So that pushes the decryption into the kernel and because the kernel was in immediately copying into the userland then it doesn't need to do it anymore because now it's decrypted and then the decrypted like can be in place the decryption can be in place in

24:29 the receive buffer itself. So it saves one copy. There's one caveat um that the Linux ciphers um are mostly not really that up to date in KTLS. Um so they might not end up using the widest CMD registers your machine has. Uh but still the mem copy saves here because we reduce one copy um offsets that um that like deficiency in a way. So we still keep this KTLS um in this land.

24:58 So what else we can improve? Well, we can improve uh more things on HTTP side but for this uh we need to do bigger changes because there's no low hanging things you can do here. But if you do bigger changes, namely writing your own HTTP client, then it pays off well because if you write your HTTP client, you control which buffers you you you give the kernel.

25:24 And if you if you do it, you can use the already pinned and page locked vortex buffer you have to accumulate and deframe on the fly in the same buffer. So once you get the decrypted bytes as a stream in the multipart get request, you keep allocating them into the same buffer you already had to be ready to be sent to CUDA. And then when the when the multipart get request is complete, you dispatch the work and you cycle these buffers.

25:48 So with one copy, you go from S3 to GPU and you have very little CPU overhead during the entire scan. And there's one more iteration of this, but you have to get rid of TCP. If you have a very expensive and tedious to set up but really performant and cool RDMA object store in your hands then you can use RDMA protocol which is different from um from TCP that it doesn't have to go to the kernel sockets layer it can bypass the CPU entirely because it works as two piece connected by some sort of fiber can be PCIe can some

26:22 sort of network connection can be in between but they can read and write from each other's memory regions. So if you have this ability then the GPU itself can dispatch the IO request on your behalf and if it does it it not only bypasses like all CPU it even bypasses the PCIe hub of your CPU because these boxes come with multiple GPUs and come with many network cards.

26:47 Most often the time each GPU has an associated network card. So the RDMA would flow flow through from the GPU into the nearest network card bypassing your CPU uh PCIe hub altogether. So it's basically free. The CPU doesn't even know what's happening. For this the network card also we need to smart to be smart enough. But most boxes that come with high-end GPUs have that network card in place already.

27:12 So with all that let's revert back to the video. So now like you can see um the there is some byte range coalesing. So we have a filter or or projection. In this case we have no filter and we project to entire we get we select star essentially all columns red green and blue. They all go through this coallesing layer and then we fetch these bytes from S3 and then we cycle the bite buffer that we read them through and that they go only one copy into the GPU and then the GPU processes them as 15 Gbits per second 4K 60 Hz

27:42 and then we stream it back to to see what's happening. But what's also cool about this because while we have it, we can change the filters and projections. We can do for example this. This is not a post projection. This is actually um a combination of a projection and a custom vortex expression. I have written this custom expression called quantise.

28:03 And I've written a GPU kernel for it which is very basic. Quantize is just gets the colors. Each pixel value is 0 to 255 and then it snaps it to three levels basically depending on who which level is closest to but you get this effect of like having only three columns. We also uh prune the blue column out today just just for fun. So it's more yellowy but also like this discretetized uh footage we have and this is importantly mostly handled by vortex orchestration.

28:33 So you don't tie this kernel at the end yourself. You just say that I want this this query and the scan orchestration takes care of appending your kernel and then dispatching work onto your kernel and then you can keep the same throughput or you can do something like this. Well, this is a bit jar to watch because the video skips all over. But the way this works is uh remember I told you when we have this file with three columns.

29:00 Well, this file has four columns and the fourth column is whether the frame has a bicycle cyclist or not. And I I've run a this object recognition like algorithm per frame. And then we now have this this has bike boolean column as the fourth column. And I'm filtering on to only the frames that has a bike in it. And the way this works is because each frame is huge.

29:21 Each frame we have is 8 megabytes. We can we a we can afford to have a zone per per frame. Uh so the in the layout level um we have a zone map that is as granular as it gets. We have one row per per frame. And so we can prune each individual frame that doesn't have a bicyclist on it. So um maybe what's more important to talk about here is all of these changes don't affect throughput.

29:50 So you can come up with an arbitrary filter, you can come up with your own your own projection, you can combine them, but you can keep the same throughput and you only read the data that you need. And remember I told you the videos were 60 frames per second. Well, because I was intentionally slowing them down. Normally, when you let them loose and then write it to write it to like a file, then we can stream 240 frames per second because I was pacing the video.

30:15 So, we can see it at 60 frames per second, but we can sustain 60 Gbits per second end to end with Vortex. So, um maybe before wrapping up, I want to talk very briefly on what you can what you can do with Vortex. um if you want to customize it. Well, remember I told you you can have your own arrays and you can have have your own layouts. A layout is responsible of IO and then logical plan and then pruning and all that stuff, but it's also responsible of placement.

30:41 There's nothing in vortex file spec that enforces you to be in the same physical file. So you can separate your segments into two. You can you can get the pruning segments in and store them in Reddus and you can get all the array segments in S3. Reddis has 100 of the latency of S3. So you get pruning nearly immediately compared to the numbers you would get with S3 and then you get the S3 throughput on on streaming the data to the end.

31:07 Or say you have this giant instance with multiple high-end GPUs but a single CPU and then you have this GPU analytics engine that you shard the data into individual GPUs and then you run a scan and then dispatch the scan into individual shards. Well, what if the query has an aggregation? Then you need to aggregate all of this results into a single GPU.

31:31 Well, you can model that aggregation as an RDMA scan because GPUs are connected to each other with an extremely fast fiber called NV link. It can sustain tremendous speed. Well, you can scan from one GPU to your central GPU to aggregate and you can scan in terabyt per second uh for this aggregation. So, uh that was vortex. Um, and I wanted to talk about like how it allows you to change your filter to change your data mix and not reprocess everything over and again.

32:06 Um, I also want to mention that you can read from S3 into GPU around network card speeds can be 60 Gbits per second, can be more if you have a stronger network card and you can read from RDMMA sources to your GPU at whatever PCIe connection you have and that can be 500 Gbits or terabs with the newest PCIe. Um, thanks for watching. Um, have there any questions?

32:36 >> Okay, we have time for questions. Thank you for the talk. Really exciting to read and listen about vortex and happy to try later. uh what I try to understand about video example you seem to not keep video in S3 but uh RGB frames that are completely separate and uh you don't apply any video codec that would uh compress data and then there would be I frames P frames that would like and you mentioned explicit that like 8 megabytes per second >> that's right per 8 megaby per frame right yeah >> so is this synthetical

33:17 example of how great is vortex or there are really use cases to fetch >> that's a really good question. Yeah and I I I forgot to mention that but yes that was a hypothetical example. Vortex is not a video codec and it was just to demonstrate that like this throughput that we can sustain end to end and like we can visualize the scan um in a video format happens to be a video um but vortex can be used like this but like uh it's mostly catered towards like analytical or training input data that you normally use.

33:45 Yeah. Is there a use for vortex with compressed codeex then? >> Yeah. Oh, you mean block compression like zl >> if you happen to have uh video stream or even uh JPEG separate but compressed data is there a value in uh passing it into vortex column to get other benefits of it? Yes, you can uh with with compressed JPEGs uh you won't be able to run arbitrary compute on them, but you can like model a census stream for a camera for example as a vortex file.

34:18 Um and we we support a lot of wide columns. So you can have like a column that is responsible of these JPEG images. But what be cooler to have is some sort of like um transparent compression. So you can for for example crop and then get only the bite ranges in that in that crop. For example, the layouts allow you to do that, but uh you need to use the right uh image compression.

34:44 >> I have a question. Are uh database engines adopting this? >> We support um data fusion, duct TB and polers. Uh we are also um in the works of getting an iceberg uh support for vortex. Uh so it's in the in the works here but ductt and data fusion is stable. You can use it um yeah today. just commented. Sorry. Thank thanks thank you for the uh talk.

35:11 Uh you just commented on the iceberg. You have some more like details when we can accept accept. >> Yeah. Um so we are working with iceberg to have a like an official vortex support. Um it's mainly on ice like iceberg is having a abstracting the file format that's backing iceberg away from iceberg itself. Um, so we are kind of iterating with them to ship this as fast as possible, but it's still in the works.

35:38 >> Yeah, I read about it like a year ago. Yeah. In your blog post or >> Yes, it's ongoing. Yeah, >> it's ongoing. >> Yeah, >> thank you. >> Nice. Thank you. I had this slightly different question like um I happened to bring this up when I was like talking to some folks at LinkedIn and then their response was does your company let you use GPUs for anything other than pure ML workloads.

36:03 Uh so like this is this seems like a great speed up for analytic data loads but like it seems companies tend to allocate their budgets for ML workloads specifically on GPUs. Yeah, I I think it's um might be even better for other like analytics use cases because uh ML workloads um you save from reprocessing your data but the speeds that that vortex can sustain the consuming model if it's a pre-training model is really slow.

36:32 So it doesn't need that much of a bandwidth but if you have some like analytics work that's happening in the GPU uh that would benefit a lot more from that higher throughput. I think >> um we still have a few minutes. Um thanks for the talk. Um two questions if I may. first um infamous sort of interoperability question I suppose and that I pick a a columnar format like parquet for instance it gives me relatively wide interrop and so in a way if I had to sort of pitch a convincing sort of narrative as to why I should

37:12 pick an alternative column format that may not have the same level of interrop what would your sort of top top arguments be and also my second question would be you talked about the sort of network penalty and the sort of CPU penalty I suppose when moving data, isn't the network penalty almost a function of how S3 and S3 APIs work? And if you had a better source, wouldn't you in a way avoid some of those network penalties?

37:35 I mean, the CPU penalty, I understand you're almost always going to get it anyway. >> Sure. Um, for interrupt, um, yeah, we have like official support for dark, data fusion and polers. Um, we are working on to integrate more. Um with iceberg I think you would kind of unlock that the the transition would be a lot easier because iceberg kind of abstracts the file from it away from you.

37:56 Um that's that's like the main idea. Um but if you want to try vortex it's it's really easy to convert from parquet. Um and like just test it yourself to see. Um the other question was about network uh limits. Um so S3 um is surprising needs a lot of tuning as well. So you can like change like the bite request even the order of those bite request uh bite range requests that you send to S3 uh change a lot of the throughput that you can get from S3 and even like you have to tune the network cards you have to bind the

38:29 right number of ENIs to the network card like that you have in the instance so on and so forth to maximize that throughput um but yes for most of the time when tuned right S3 itself would not be your bottleneck in this pipeline Hi. Um, most of that pruning stuff that you showed seemed like it was reliant on the data being sorted in the first place. Uh, biggest problem we've ended up having with a lot of this stuff is that lots of different users need the same data sorted in 10 different ways.

39:02 Could you use something like the layouts to basically keep all the metadata that you need so that you could effectively have the same data sorted 10 different ways for people to read? I'm not sure it would be as as efficient. you'd have to pick a most efficient option but >> yeah yeah yeah be an option we mostly delegate sorting itself to the compute engine that we're using so data fusion or or like DB that would handle it uh but yeah I think it's an option >> hi I have a question about the uh metadata we store in the

39:31 layout so we talked about min max of a particular field like how do we decide what metadata to expose like in the case of age like the query was designed and the minmax was there. But yeah, in general, how do you how do you write vortex files knowing um what might make it better? >> Yeah, sure. Um so this um actually um is is um transparent to you. So you don't get to pick like which the zone maps like will end up being in the file.

39:59 But the statistics that we support, I showed just a subset of them. But you can like anything that you can do um to calculate before end as statistics when writing that you can use to prune afterwards. We will sort it like null count for example like min max are examples you can have um you can have even cardality estimation that you have in your file.

40:20 So you can help up with like aggregations and stuff. So anything that would can help with pruning we we try to integrate it and adding a new one um is is really easy in watch if you want. I had one question like in one of your slides really early on you had the dictionary in uh encoding. >> Yeah. >> The three colors. >> Yeah. >> And I think you said there's like 33 uh four bytes for the data.

40:43 But then you didn't count the bytes in the dictionary itself, right? >> That's right. Just the just the overhead was around 30 bits or so with Yeah. So this is like 30 bits and that's plus this. Got it. >> But you can compress this further uh if you have the right encoding. Yeah. And maybe to finish on the S3 question, um actually like maybe counterintuitively, but S3 uh bandwidth can surpass the NVME bandwidth if you tune it right.

41:15 So for example, in in those scans that I was showing uh you can you can scan from S3 fast faster than you can scan from your local NVME. The initial latency is higher but the throughput is also higher from S3. Um maybe I can ask a question like so I remember when Park first came out. >> Yeah. >> Right. This was when like the big data world was trying to compete with the established uh MPP databases.

41:49 And one of the first things they did was they said hey like row major is not great. Let's go to this column format. The next thing they said is well we're using JDBC we're taking what was columner and going back to row major and then we're not getting any of the benefit with pipelining this through memory through all the m layers through to registers that's where arrow came in >> right so how does like arrow figure into this >> yeah it's a good question um so all the decompressed um vortex like we call them the

42:22 canonical encodings for each logical type we have a canonical encoding they are one to one with arrow and there are zero copy with like from and to arrow. So you can that's how we plug into data vision as well. Data vision needs arrow batches and we just decompress vortex arrays get the arrow batches and then plug let forward that. But if you have for example duct DB, duct DB has multiple types of vectors.

42:41 There's a constant vector and we also have a constant array. So it's already a constant array. There's no need to decompress and then re-encode it by duct DB as a constant array. So we kind of pipe that through. Um okay. >> So let's say I take a parket file and convert to vortex. >> Yeah. Um, you said that it's I think you had stats on how much faster.

43:01 >> Yeah, I I have numbers. Yeah, for >> I think it was before if you like in the beginning slides, right? You had like one something like 30x faster or something. >> Yes, I have numbers for GPU scans, but I also have numbers for CPU scans and rights and random. >> Okay, but not for you. >> So, um, this was telling that the GPU scans were 30 times faster, 100 times faster for random access.

43:21 The CPU scans around 20 times faster. CPU writes writing a parket file is run around five times faster than par in that ballpark. >> So it's it's sort of I think to the other question which is like you know avoid >> and then it was a huge amount of data but they converted it to park because it made sense. >> Yeah. >> So I guess it would also make sense eventually to move all park to vortex >> if it's faster on CPU and GPU.

43:50 Is there any size difference in the actual file? >> Size difference? Yes, the compression ratio is very similar to park. It's like 10% more, 10% less in that same ballpark. >> Okay. So there's like the like uh compression of like park and then this compression like park zstd. Yes. >> Which is even further. So like you're are you comparing park to vortex or >> std park to vortex?

44:17 >> Oh par zstd to vortex. You don't need to apply zstd on vortex. >> No. And would there be any benefit to doing that or would you basically not get much compression? >> So we have we have so the layouts um they also have writers and readers. We have a custom layout writer that you can use called the compact writer that does the standard for some data um but for example uses PP pcodc for floats.

44:38 So if you want a small even smaller file you can tune that in but then you would lose the computation advantages that you would get. Um so it's up to you and it's vortex is completely like configurable. You can have any writer reader and combine them as you wish. >> Okay. Another slight different question. Sorry if I'm like this is V1 or V sub one. >> It's V1.

45:02 >> V1 stable as of 6 months ago I think. >> What what is uh and you're a committer maintainer on this project? >> Yeah. >> So like if we look at protobuff it's protobuff 3 over 20 something years right. So every seven years something they changed. >> Yeah. >> It became incompatible with previous. >> Yeah. So what is sort of your view on vortex like so one thing is like if we look at GPUs they have this uh virtualization layer >> and you you code to that virtualization layer the GPUs change very often but the this the

45:33 CUDA coding doesn't change at all >> right uh how do you view your format this format what do you have to keep up with because if you're coding to that layer versus to the GPU hardware itself then you're kind of standard But you have to code to the you have to change this to make the most use of the hardware itself. >> Yeah. >> And the hardware changes very fast.

45:56 >> Yeah. >> So how would you see the backward compatibility of vortex over time? >> Yeah. Um so the the the file spec is standard and it's uh intentionally minimal. Uh so we don't need to change the file spec to change most of the things that you would need a file change for for other formats. For example, like parquet spent around seven years incorporating F-16s.

46:18 Uh because they have to decide the the committee has to decide. A lot of big corporations behind it have to agree on one form or another or ALP for example is a cool floatingoint compression is been open for like a couple of years now and it's like nowhere near to be merged. But for vortex, it's just a crate. You just write a new rust crate and then that that's your new layout.

46:38 So for um so the the the core part of the file would not break because it's very simple. What you can have is you can have some new encodings that other people have written that you don't have. So you cannot read their files but you can just add the library and then read them. Um or yeah >> you support non rust serties. >> Yes, we have an FFI binding.

46:56 Uh we even have a Java uh Java binding. So yeah, any any languages you can plug in with vortex. >> All right. Sorry I totally manipulate. I totally took over this. you gave example with arrays and dictionaries. Um it's like in the beginning of slides. Uh but ultimately it's still you keep nested data structures inside a single column. you optimize them inside a column but you don't you still don't attempt to expand them on a sort of shredding pair of nested structure spell columns >> to achieve like >> better

47:40 compression there. So my question is is there a step improvement in vortex uh versus parket for maps or JSON this kind of data structures? >> Yeah. Um so we currently don't have that in the core u but it can be easily added as a as a custom layout because shredding is just like a new layout that you would read you would you would write yourself and uh it would treat the JSON blobs as binary and then extract that and then store them store the fields that you want as a separate column.

48:07 So it's an ext great extension point. It's just not in core yet. >> Is it yet not in core or it's sort of by design not what you're trying to >> Well, we we we think the compression ratio and then the query patterns that for example tpcds and clickbench they like this compression kind of like works a lot like works well so far. But if we need like some use case that would benefit greatly from shedding then like it's um not that difficult to integrate in with a new new type of layout essentially.

48:43 >> Um I have a question regarding uh okay last question sorry >> uh regarding consistency and error correction about this vortex file. um which mechanism are you using and uh how when we read the file you can identify a file is corrupted and so you need to redownload it or recreate it. >> Good question. Um currently by default um we rely on the network stack I think um but you can add like a Merkel trees or like any sort of parity on the photo as metadata and you can get your layout to check them uh before before

49:16 dializing the files if you want. Okay. Um, thank you very much. >> Thank you.