Open Table Formats in the Wild: From Parquet to Delta Lake and Back

Description

Open Table Formats (OTF) such as Hudi, Iceberg and Delta Lake have disruptively changed the data engineering landscape in recent years. While the Parquet file format has evolved as the de-facto standard for open, interoperable columnar storage for analyical workloads, it lacked first class support for critical features such as ACID compliance, incremental processing, flexible schema & partioning evolution and scalable meta data management. This led to increased development and maintenance efforts while building idempotent and failure tolerant data pipelines that often resulted in custom frameworks. OTFs solve all of these issues via providing a sophisticated meta data layer and improved maintenance capabilities on top of Parquet.

Driven by the promises of OTFs, we intended to replace our own bronze-read-only Parquet-based storage layer with Delta Lake. In theory, this should have improved performance, reduced maintenanced and provided more flexibility. However, we've stumbled upon several issues:

  1. drastic performance issues with Liquid Clustering during incremental processing
  2. inmature interoperability in the python and cloud-based ecosystem (DuckDB, Pandas, Polars, Athena, Snowflake)
  3. maintaining logical session-boundaries during incremental processing

While the first two issues are solvable in foreseeable future, the last one is specific to our requirements and does not overlap with design decisions made for incremental processing in Delta Lake. Taken together, these points ultimately led us to go back to relying on Parquet again.

Targeted Audience

This talk is mainly intended for an intermediate data engineering audience but is well suited for interested beginners, too. The content of this talk is relevant for all architects and data engineers being responsible for storing and managing data for analytical workloads.

Key takeaways

  • What problems do OTFs solve?
  • How do OTFs contribute to an open, composable data stack?
  • Is there a predominant Open Table Format?
  • How does Delta Lake conceptionally work?
  • What are concrete real-world advantages of Delta Lake in contrast to "plain" Parquet?
  • What is the "small files" problem and how does Liquid Clustering help?
  • How is the current state of interoperability with Delta Lake?

Talk Outline

  • Introduction (5 min)
  • OTFs in comparison (5 min)
  • Delta Lake Internals (10 min)
  • Use Case Requirements (5 min)
  • Benchmarks & Results (10 min)
  • Conclusion and Outlook (5 min)
  • Questions (5 min)

This session took place in track Data Handling & Engineering and was classified suitable for intermediate domain / novice python by the speaker.

Transcript (auto)

Auto-generated from the recording utilizing Open-Source AI. Speaker labels (Speaker 1, Speaker 2) reflect diarization, not identity. Timestamps refer to the recording.

Speaker 1 [00:07]

We are also very well welcome from my side and I'm excited to be here on PyCon, I'm excited to be here on stage and I'm very happy to talk about Open Table Formats today because I think they represent a technology that has changed the data engineering landscape in recent years to the better. But before we get into the very detail, let me quickly outline the agenda of the talk. First of all, I want to start with a brief moment of thinking about why data is important at all. why are we here? Why care about storing data? Isn't that a boring topic? Turns out it's actually not. And then I would like to invite you to two histories. One history lasting 100 years, a century of development of printing machines, and a related history lasting only 10 years, which describes the very first big data platform that we used in my company. And this already introduces us to Apache Parquet. We will then deep dive on Apache Parquet and see and understand how Apache Parquet works. But we will also showcase the issues that Apache Parquet didn't solve, which essentially gave birth to the existence of OpenTable formats such as Delta Lake, Iceberg, Hoodie, or also Apache Payment. And then understanding all the promises that OpenTable formats have made, we want to put a specific OpenTable format, Delta Lake in this case, at a test. because we thought, well, why not replace our formerly Paki-based Bronx data layer with an OpenTable format to improve performance, to reduce maintenance, and to gain more flexibility. And it turns out this isn't as easy as we thought it would be, and there's a twist, and it will be very interesting. And in the end, I would like to take a very short moment to talk about the current state of table formats. As I mentioned earlier, there's Delta Lake, there's Hoodie, there's Iceberg, There are also other formats, such as Apache Payment from the Flink community, and I want to give you my personal opinion on it, since it's currently a very hot topic. So, let's kick it off. A 10,000-foot view, it all begins with data, but why? Why are we here? Why is PyData so popular at the moment? It has been popular also in recent years, but just this morning we heard more than 2,000 participants. Well, typically before I start to prepare a talk or a presentation, I ask myself, why is this relevant to me or to you as the audience? Why is it relevant in general? And it kind of struck me after a while that data is at the very core, at the very heart of scientific progress in human advancements. While this may sound or seem obvious once you see it, it is not obvious right from the beginning. When you think about your personal life, when you think about your university studies, You come up with some ideas, you come up with some theories. We as humans, we want to understand our environment, we want to improve our quality of life. And how do we do this? We do this by understanding relationships in the outside and internal world. And coming up with ideas and theories, we somehow have to falsify or validate them. So we need to collect some sort of external criteria that improves our hypothesis to be right or wrong. And this is done by data. We observe and we collect data. And sometimes we have an object of research which is fairly simple. For example, imagine measuring the weight of a stone. Fairly easy. One data point that I can memorize in memory. Or the size of a person, also fairly easy. But there are other examples which are more sophisticated, which are more complex. For example, I have a background in psychology, and in psychology we study human behavior and cognition. And we all know that we are humans, we are so vastly different, There's so much variance, there's so much random error involved, and you need very dedicated questionnaires, dedicated measurement devices, and dedicated statistical methods to reveal a pattern to gain an insight. And it's all about having data and being data-driven to gain insight to progress with the scientific research. While this is rather abstract for now, let's take a look at a very clear example, our first 100 years journey. And this is about printing machines. And this is not by accident. It's because I'm working for a company, Heidelberger Druckmaschinen. Heidelberg is the name of a town which is roughly 60 kilometers south of Darmstadt. And Druckmaschinen is the German word for printing machines. And this year we are celebrating our 175th birthday. And I also learned something new during this anniversary, that we didn't start off with producing printing machines, but rather we started off producing belts. For example, one of the original belts in the Cologne Dome was from Heidelberger Druckmaschine back in the day. But it's not about the history of Heidelberg, it's about the history of those printing machines and how data has affected those machines and their efficiency. On the left side, you can see an original Heidelberger TIGL, which was developed roughly a hundred years ago. And on the right side, you can see the most recent machine, it's called a Speedmaster. And let's compare these two machines on three high-level KPIs. The first high-level KPI is the number of sheets printed per hour. The Tegel on the left could already print 3,000 sheets an hour, which I think is impressive for a machine being 100 years old. The Speedmaster masters the speed, as the name suggests, prints 21,000 sheets an hour. So this is like a seven times difference. Now let's take a look at the second high-level KPI, the number of colors applied. The Tegel on the left side could apply one color at a time, whereas the Speedmaster on the right side, you can see those units, and there are eight of them. There are four units, four colors applied on one side of the sheet, then the sheet gets turned around, and then another four colors applied on the other side of the sheet. And this is a very standard configuration. You require four colors in order to cover the entire color spectrum. So this is an eight times difference. And now the last highlight of the KPI is about the sheet size. With the TIGL, you could roughly print an A3, like this size. With the Speedmaster, you can print A0, like four times the size of the TIGL. And if you multiply all this together, you end up with a 244 times increase in efficiency between those two machines during a period of 100 years. You could also put it differently and say every 10 to 15 years, we double the efficiency. And this did not happen by accident, but it happened because we used data, we used theories. And these machines, you can guess on the left side, the data ingestion or the data tracking was a rather manual process, whereas on the right side, there's lots of sensors inbuilt to automatically collect data and store data. You can think of low-level IoT sensors such as temperatures, currents, speeds, error failure messages, consumption messages, but also high-level KPIs such as in OEE or the number of good sheets or wasted sheets, make ready time, it's everything you can think of. And with the ever-increasing amount of data and the number of increasing machines connected, we are facing an issue on how to store and process all of this data. And this brings us to our next history that we are going to cover, the 10 years journey of the very first big data platform that we used at Heidelberger Druckmaschinen. And back in the day, 2015, ten years ago, Apache Hadoop was kind of state-of-the-art. So Hadoop used an HDFS distributed file system to store all the data, and the data was already stored in Parkey format. There was a very wise man who decided to use Parkey because this is something we still use up until today. We then also added Hive to the mix and improved on having a logical table instead of a file-based access. We also use Cassandra and Kudu, but it's not so important. We used Spark as a batch processing engine, and we also used Presto for interactive queries. And fun side fact, with Spark, we still have Spark jobs in production that run in Spark version 1.5. They've been written with Java-based API using the RDDs instead of declarative data frame APIs, and they even included a custom domain-specific language from our external partner. And Toby, one of the guys who wrote those jobs, is also here today unfortunately not with heidelberg anymore but it has been a torture for him and has been a torture for me to maintain these ones and i'm very happy if we can shut down this platform by june this year um why do we shut it down well first of all um with hadoop you have a couple to start and compute you can't scale your storage independently of compute which resulted resulted in a scenario where we always moved the data that is older than two years to an A three bucket to an external bucket that was not stored on the hdfs because we couldn't scale it independently um also um you can't just increase compute and storage just like this as you can do on aws for example or in google cloud or on azure um in our case this means on mondays when most of our customers that their shifts we get lots of incoming data then there's a high peak of processing and And there's a huge backlog that builds up. In contrast, during the weekends, there's no incoming data. The entire platform just idles and wastes money. And there's also another reason why this failed. We as a team, our internal team, also failed to maintain up to the standard and up-to-date framework, such as we didn't have a proper, we still don't, up until today, we don't have a proper orchestration framework, such as Airflow or Dextro or Prefect, to model interdependencies between Spark jobs. We don't have infrastructure as code framework to reproducible provision our Spark jobs. This is rather a manual process. And there are many other things that didn't go that well. And well, what we did is probably what most did. We went to the cloud. So we went with AWS in our case. And the cloud world is so full of promises, too. For example, it's so easy just to open up an AWS account, create in a three bucket, and just store tons of data there. And it's cheap. decide, well, I'm not accessing the data that frequent, so I just choose a different storage tier. And then I can spin up whatever compute I like. If I have a very compute-intense job, then I use compute-optimized EC2 instances. If I have a rather memory-intense job, then I just choose some memory-optimized instances, or I can use some GPU instances. Not possible with the traditional Hadoop platform, no way. But, and there's a big but to it, it's not that easy to just move from Hadoop to the Because if you want to orchestrate all those different services that you require for an enterprise platform, this can't be done by a data engineer. It can't be done with a single team. You have to integrate storage, compute, orchestration, monitoring, alerting. Everything has to be testable, reproducible. You need something like authentication, authorization, user provisioning. This is like a whole enterprise project that cannot be done by an average-sized company like Heidelberger Druckmaschinen. And that's why, you know them, other companies have done such things. They have built these integrated data platforms, such as Snowflake and Databricks. And at Heidelberger Druckmaschinen, we actually use both. We use Snowflake and Databricks, and this decision has been made two years ago. And why is that the case? Why not just use one of them? Well, back in the day, two years ago, Snowflake had a clear edge in regard to warehousing and BI, whereas Databricks had a clear edge when it came to more complex data pipelines with a native Python interface, which require proper testing, introspection, and also the machine learning frameworks integrated in the Databricks platform were way better than Snowflake once. Nowadays, they consolidate more and more. They get more and more closer regarding their functionality. Snowflake now offers native Python notebooks. Snowflake now even offers a managed Docker runtime. How crazy is this? Whereas Databricks, yeah, they started with Databricks SQL, but they are now also competing in the warehouse space. And Databricks also now offers like a complete serverless setup. You don't even have to touch any AWS resources anymore if you want to run Databricks. So they're getting closer and closer together. But besides Databricks and Snowflake, well, if you're interested, come to our booth. We love to talk about this. But besides these two, Parquet persisted. We still use Parquet today. So let's say hello to Parquet again and understand why Parquet is superior to CSV or JSON for analytical workloads. Why is that the case? Why is JSON and CSV guys looking jealously at the Parquet hugging the data here? So for one, Parquet supports strict types. This is not the case with CSV. It can be the case with JSON, if you use JSON schema, for example. But strict types, they offer a huge benefit. First of all, they convey meaning. And secondly, they remove any serialization and deserialization issues. When I started with Heidelberger Druckmaschine as a data scientist eight years ago, I needed to write tests for my Spark data pipelines. And I needed a format that can also communicate with non-technical users. So I couldn't say, use Parkey, for example, back in the day. But rather, okay, here's Excel. Everyone can use Excel. Please provide your test data in Excel. But then I had to get this Excel data into Spark. And this was a hell. It was even before Pandas 1.0. So this meant using Pandas to load Excel. You get a hell of type conversions of timestamps, for example, or even strings. And then I had to convert Pandas data types into Spark data types. This was also a hell. And then, to do the testing, Spark back in the day did not offer a native way to compare data frame equality, which is kind of surprising. Pandas did have a possibility for this, but Spark not. So I had to convert it from a Spark data frame back to Pandas to finally test data frame equality. And there was so much boilerplate code just involved to test. Anyways, the next benefit, efficient encodings. Parquet offers a whole lot of built-in encodings such as dictionary encodings bitcasting bitpacking bitpacking run length encoding and you can even chain them together. For example if you have a string column why store these expensive UTF-8 strings or characters rather just map them to numbers use a dictionary encoding and let the CPU crunch with numbers for your group by aggregate query And then, for example, if you have a string that has lots of consecutive repetitions, then you can add a run length encoding on top of the dictionary encoding, saving more space and making your CPUs run faster, typically. And the last one, Parquet files are splittable, which we will also see shortly. That means if you have a multi-core setup, you can process a Parquet file in parallel Instead of, for example, having a JSON, try to split a JSON in part. Or if you have a GZEB CSV, try to do this. It doesn't work. But I think the most striking feature that Parquet added here was the possibility to skip irrelevant data. And now I want you and me to assume a role of a Parquet writer application. On the left side, we have an in-memory data frame. Three columns and six rows. And we will now follow the Parquet writer specification, or the Parquet memory layout. So the first thing we do is that we partition our data horizontally. We create row groups, groups of rows. Next, we partition vertically. That means we take column chunks out of these row groups, and we store them sequentially on disk. And this is important. Typically, we say Parquet is a columnar-oriented file format. And yes, it's column-oriented, but also we split horizontally here. We have row groups, and we have those column chunks. And why is that the case? We are talking about analytical workloads. We're interested not just in a single row, but rather we're interested in entire columns spanning a long history. There are no CRUD operations involved in our analytical workloads. But rather, you have a filter condition on one column and an aggregate clause on the other column. And hence you want to access entire columns, but also you want to have these columns to be to be co-located If columns are stored somewhere else then still you have to read large chunks of your data But in this case we have kind of a sweet spot. We say we have columnar layout, but also row groups split it horizontally So we chunk these columns A B and C and while writing out the data we collect statistics about these columns such as the min and max values and we put these statistics into the footer along with byte ranges that allow us to specifically only load this very part of the data instead of loading the entire file then we do the same with the second row group and once this is done we have yet another footer for the entire parkey file and this footer contains information about the row groups and now let's put this to a test now we are not We are not the parkey writer application anymore, but rather we are the parkey reader application. So we have an on-disk format, on-disk storage, and now we want to read data from this. So what do we do first? We want to get column A, where B greater equal three in this case. So first what we do, we only read the footer, nothing else, just the footer. Given the footer, we know that row group one's largest value for B is two. So we can completely ignore row group one. We never have to touch it. We can completely skip it. And this is called predicate pushdown. There's a filter condition and giving this meta information that we have by footer, we can completely skip scanning row group one. And then in the second case, well, we need A and we need B for filtering. Then we have the projection pushdown and both works with Parquet. This is not possible with CSV. There is no concept of metadata in CSVs. And that's a striking feature. And this is also a feature or an idea that is applied once again by the OpenTable format. when dealing with lots of Parquet files. Okay, so I hope I have you convinced that Parquet is good. But there are some downsides of Parquet. For example, the Parquet specification does not contain anything about asset compliance. So if you have multiple write applications, transactions won't work. Performance and scalability. Once you have a lot of data, Parquet files themselves alone tend to be slow because of all this footer information. You have to query so much footer files in order to know where to look for your data. And also when it comes to schema and partitioning evolution, Parquet is not very flexible. And this brings us to OpenTable formats. So on the left side, we have Iceberg as this Iceberg. On the right side, we have Delta Lake being this other animal. And in the middle, they are hugging Parquet because they are still using Parquet underneath. But now they're adding yet another meta layer on top, yet another footer on top to make it more efficient to work with Parquet. And the yellow box is all of the things that Parquet brought us. And the blue box are all the additions made by OpenTable formats, such as schema and partitioning evolution, asset compliance, stream and batch processing. This is a huge thing. Incrementally process your data. And there's no need to have some external application to track what data has been processed or not in the future. if you just rely on those OpenTable formats because they have an internal versioning. Native maintenance possibilities, as we will see later, time travel mutability, you name it. And how is that the case that now Iceberg and Delta Lake can manage all of these Paki files, all of these little Paki beavers in this case? And as I said before, they just add yet another layer of abstraction of metadata. and let's have a look for data like specifically in this case what does the metadata of data like contain it's like in transaction log because you can see here on the left side a different path each path corresponds to a single per key file and then your transaction or transaction log is like a binary log of your typical transactional database tells you well per key file zero was added and then we have some min max values we also have from some partitioning information but you can also delete files and you can also change the schema with this DDL data definition language and then you can version it and having different versions means well if I have a reader application it remembers the last time I came here I only saw the data of version one now I have to read all everything that has been added afterwards so you have this incremental processing allowed time travel is now possible and you keep all of this meta information in this single place. With Iceberg, it's a bit different. The idea is the same, but the structure is somewhat different. I won't focus here because I don't have that much time on it, so we skip this for now, and let's just see examples of transactions and time travel. One important note here, if you want to support transactions, you need to have some sort of external locking mechanism. For example, with Delta Lake, typically nowadays you would use a catalog that does it for you but you can also rely on a DynamoDB table to do the locking for you. It requires some sort of atomic swap in order to make this new version be the newest one. Let's see a typical scan without metadata. So we have those 12 parquet files and we are only interested in reading the blue one. Without metadata what our engines do well yes I scan all of these parquet files and their footer information because I just don't know where the data is that I'm looking for. Well, however, if I have metadata available, such as for Iceberg and Delta Lake, what I do well, first I read the metadata and then I know where to go. Then I know, oh it's there and I only have to look there. Two API calls. And if you scale up then it becomes very obvious that you have a huge benefit if you have metadata available. And the next one is not clearly related to Paki but more with Hive partitioning. Let's assume we have one partition and there's lots and lots and lots of small parkey files there. First of all, you get a lot of metadata overhead, metadata overhead of the parkey files themselves. But also, think about the HDFS name node, for example. The HDFS name node is responsible for keeping a reference of each parkey of each file name to their corresponding storage location on the data nodes. And this HDFS name node has an in-memory presentation of this reference. And it's not unlikely that you run into an out-of-memory error, because you just have too many small files. The same applies to our Delta Lake and Iceberg index, because they also have to keep track of all of these files. So this is not very efficient. And also, when we think about cloud storage, cloud storage for analytical workloads is very efficient if you have large sequential reads. But in this case, you have many small reads. This is not very efficient. There's lots of overhead involved. So what you can do, do compaction. We just move all of these tiny parquet files into one large parquet file. That's the way it has been done in the past. But now let's assume there's a second partition, but this one only has less data. And then we also compact it, but then we have a different size of parquet files here. not really our target size we're interested in and this is a diagram taken from databricks and which showcases the difference between is the fixed partitioning scheme used by hive which results in skewed file sizes and the difference on how liquid clustering as you soon see improves on that in this case yes we hit the target size for some partitions but some partitions are completely empty and other ones the newest ones there are plenty of those small files have not that have been yet compacted and if you scale out you still may run into a problem there are lots of small files and partitions are not evenly balanced in this case you can apply liquid clustering it allows you to have flexible partitions and equal file sizes and as we see for example in this case in the us in february 2023 we now just create more partitions instead of having one partition with lots of data and this is what i meant earlier with maintenance improvements Liquid clustering has a specification, is not implemented in every writer application, but when you use it, it can improve your quality of life as a developer. So, promises made. Reduced maintenance, improved performance, more flexibility. Let's go for it. And for us, we wanted to substitute our Parquet-based read-only Bronx layer with Delta Lake. So, put it to a test. Under investigation, liquid clustering. Because it seemed so attractive to us to use liquid clustering since our per-key files, we also suffer from a small file problem. Our per-key files typically have a size of 4 to 5 megabytes. That's not enough. So what we did, we had a setup in which we used 1.3 terabyte of session files. You can imagine a session file containing a machine number, a timestamp, and then some arbitrary information. That's the most important part for you to know. And then we have on the left side classical hive partitioning. so fixed partitions, and on the right side, we have liquid clustering. And then there are also two important distinctions here. We did this benchmark while processing the entire 1.3 terabyte at once to create the Delta Lake table with liquid clustering, which does not really apply in our day-to-day work, because typically we process roughly 50 gigabytes per day, and then we add it incrementally to a Delta Lake table, and then it has to be clustered. that's you will see why this is important so let's take a look at the very first example so we read 1.3 terabyte of data and write it to data lake with liquid clustering or we have it with hive partitioning let's look at the chart on the y-axis you see the total number of api calls required to fulfill a query by machine and timestamp so we create just a random machine and then also provide an additional filter on the time span and then you see the the is it orange for you it's orange the orange bars they correspond to hive partitioning and the blue bars they correspond to liquid clustering the lower the better and we can see as expected given that with metadata information available flexible partitioning liquid clustering is clearly the winner. Another KPI, the total duration taken by those API calls, and yet again, we see liquid clustering is clearly better. So is this really the better life? Well, we thought initially it was like, yeah, hooray, I'm talking to my manager, we can improve here, we can save a lot of money. And then I did another test using the incremental ingestion. So in this case, not reading 1.3 terabyte at once, but only 50 gigabytes, and then adding it to a Delta Lake table and clustering incrementally, and all of the sudden, the effect almost turned around. So liquid clustering got worse, even worse than hive partitioning. And in this case, with the total duration, it's actually worse than hive partitioning, and that was a huge surprise for us, because the promise of liquid clustering is, well, use liquid clustering. You don't have to care about clustering anymore. We do it for you in the background. So obviously, it fails to properly rebalance the partitions with default parameters. So we reached out to the Databricks support, running this on a Databricks runtime, and they told us, well, there is a configuration property that is not publicly documented, but you can use to tweak how liquid clustering behaves. All right. So we did this, and yes, it worked then. But it doesn't hold up to the promise that you just relieve the developer or the data engineer from manually clustering and manually observing performance of your table. And what's even worse, in my opinion, is currently there is no way to access the clustering health of the table. And when should I use this skewness threshold, when should I apply it? And they're saying themselves, well, without a method to monitor it, it's difficult to determine when to exactly use this. So this kind of left me unsatisfied. And just to make a comparison here, Snowflake proprietary format, I'm not a big fan of proprietary formats. about open table formats for their internal proprietary format if you use clustering it's kind of a similar thing that liquid clustering does at least they have a method that you can use to see whether the table is clustered well or not given certain columns okay um the second thing on an investigation the idea of interoperability how well is delta lake supported out there in the python and cloud ecosystem the idea is right once read everywhere naive defaults what we did We used Databricks Runtime 15.4 LTS and stored it on E3, and then we provided access to the E3 bucket to all of these different engines and wanted to see if they can just read and operate with the Delta Lake table. So, we used DuckDB, and you can see, well, DuckDB, at least it knows the schema, but empty data. This was not successful. We used polars, and oh, there's an error, an exception, delta protocol error, certain features are not supported. Okay, column mapping, V2 checkpoint, deletion vectors, I didn't know about them before. Then using pandas, the same error message, that's because they're using the same Rust kernel underneath. With Amazon Athena, the clue error just spit out an internal error. I don't know. With Snowflake, empty table. So this was not really satisfying. So could you say, write once, read nowhere, and say in defaults? Turns out, Delta Lake has table features. I'm not here to bash Delta Lake or Databricks, don't get me wrong, I'm just saying you have to get your head around and understand what's going on. On Databricks, you get the newest table features available and they are activated by default. But if you want to use open source engines to read the data, most of them do not support them yet. And sometimes, at least for me as a user, it was confusing to see which one is supported, which one isn't. And also, you find this kind of strange that on this slide, this parameter here, it says spark.databricks. And even the open source Delta Spark implementation has spark.databricks. Why is databricks there? I thought it's Delta Lake open source. So this is kind of confusing. Anyway, so tweaking this Delta Lake table, removing certain features, finally brought us where we wanted to be. Now we have a DuckDB reading the data successfully. Ah, with Polars, somehow there is still this timestamp with native time zones, so not a time zone aware timestamp. This is not supported yet with Polars and Pandas, but Amazon Athena finally worked, and also Snowflake worked. So as a result, mind the difference. You have to know which features you use. And the last one, this is rather special in our case. We wanted to do incremental processing. We had a manual framework for this in the past, and we wanted to ditch it because development work, maintenance work, and so on. And there's one important thing to understand in our case is that when a printing machine boots up, and until it runs, until it shuts down, a single session file is created. A single parkey file is created. And to make developer life more easy, we wanted to just reason about entire sessions instead of reasoning about partly available sessions. So instead of having a streaming framework, we wanted to have batch processing, because batch processing is both more efficient and easier to write and easier to test than streaming. So we wanted to preserve these session boundaries. And what you can do on Databricks, you can use autoloader. And autoloader works perfectly with file semantics, incrementally loading entire files. This is very similar to what Snowpipe does on Snowflake. And if you use autoloader, it's just fine. you can stream entire files, entire sessions into your Spark streaming context and then apply a batch processing procedure on top of it. Perfect. But we thought, well, maybe we can just go without autoload at all, just rely on pure Delta Lake open table specification because there is incremental support available. But it turns out Delta Lake has table semantics. The session boundaries that we had via entire Parquee files, they're not preserved. And this is by design. There is no way around. We also reached out to Databricks support, and they said, well, your use case is nowhere near the design of Delta Lake or OpenTable formats in general. So this was also kind of not really satisfying, but it's like a homemade issue that we have. So summing up, we wanted to go with Delta Lake because of all of these promises, but then realized, well, liquid clustering didn't hold up the promise. Interoperability was still challenging. and also our specification to preserve session boundaries is very specific that can't be really forced into the design of OpenTable formats. Iceberg and Hudi also won't solve this for us. Okay, so that's so far for Delta Lake and Parquet. Now a few words about Iceberg here on the left, our left contender, and Delta Lake on the right. I left out Hudi, because Hudi is more focused on streaming workloads, even though we can do patch processing, and also left out Apache Payment, because this is rather new, coming from the Flink community, which is also more focused on streaming workloads. Let's have an example of the main contributors to Delta Lake and Iceberg. I did this roughly two weeks ago, so it's representative. And on the left side, we see all of the contributors to Delta Lake, and even the first one is someone from Databricks. And it's almost exclusive people from Databricks. Whereas if you look at Iceberg on the right side, you can see many different names. And a lot of these Gmail addresses or names, they are from Dreamio, they are from Starburst, they are from various companies. And you can see that the community engaged with Apache Iceberg is more diverse, being supported from many different companies, whereas with Databricks, it's almost only Databricks at the moment. Secondly, if you think about the overall community engagement with these two table formats, just two weeks ago, there was the Iceberg Summit for the second time, one day in person and one day remote. It's being supported by many different companies, from Dreamio, Snowflake, AWS, Microsoft, also Databricks. And there's like a whole own summit dedicated to Iceberg. And if you look something for Delta Lake, well, the only thing you will find is the Databricks AI Summit in June this year. And those are the five talks that popped out after a search for Delta Lake. That's not a whole lot. So it seems from the community perspective that Iceberg has an edge. And there are also other cloud developments developments such as F3 tables being announced by AWS natively supporting Iceberg. And yet another note, Databricks bought Tabular last year and Tabular was one of the main pushing forces behind Iceberg and we wanted to enable true interoperability with our Databricks runtime and iceberg so one thing you want to do is having an external data catalog connect to it and write iceberg format natively with the databricks spark runtime but it turns out that even though vanilla open source spark is able to do this the spark runtime on databricks isn't and at least for me this was disappointing and to be honest was just disappointing you can use vanilla open source Spark it can do so but the Databricks runtime and managed enterprises Databricks Spark runtime can't it's on the roadmap, you can read Iceberg data from external catalogs if you integrate them, but you can't write data and that's disappointing yeah so in the end, I think it doesn't matter if it's Iceberg Delta Lake hoodie or whatever in the end open table formats. They make our data engineering life better Just remember a few years ago if you use snowflake redshift to the bakery Whatever your data our data is locked behind proprietary data formats And it's only accessible via proprietary cloud engines owned by Wenders so the entire Ecosystem is changing now. We're using open table formats such as iceberg or Delta Lake And we're not relying on any proprietary formats anymore. And it's our data, and we can access it with any engine we want. We are not forced to use Snowflake Compute Engine or Databricks Spark Runtime. We can use DuckDB, we can use Polar, Pandas, whatever. DataFusion, also to name it. And I think that's a great thing. That's a great future ahead. That's also the reason why this iceberg animal and the data lake animal, they have the same logo right here. Doesn't matter if it's iceberg or data, who cares? just have open table formats making our data available removing or reducing the lock-in that formerly cloud proprietary data warehouses such as Snowflake and Redshift had. This is my closing statement for today and thank you for attending. I'm curious for your questions.

Speaker 2 [39:00]

So, thank you for the talk, Franz. The first question is not actually a question. It's just the audience that wanted to say that they loved your visuals that you introduced into the presentation. So, the first real question is if you could elaborate on DynamoDB for locking files as opposed to using Polar's Unity or any other catalogs. Is there any conscious trade-off?

Speaker 1 [39:25]

Well, if you want to use DynamoDB just for locking you just use it for locking. I think don't use it It's just there that you can use it to have to enable transactions to enable multiple processes to write to a single DataLag file. It's very similar if you for example have used Terraform You also need some sort of locking mechanism that you don't override or corrupt these Terraform state is the same idea You just need some external criteria to tell you whether you can write to this metadata or not. It's just a lock that you use. And in this case, it's DynamoDB. But I would always rely on a true data catalog such as Polaris, for example.

Speaker 2 [40:05]

Would you say there's something like liquid clustering with iceberg tables?

Speaker 1 [40:10]

It's an interesting question. I also thought about it and I think currently there is no implementation for this, but there could be an implementation given that Iceberg also allows partitioning evolution and since partitions are just metadata in the hierarchical metadata file of Iceberg, there is no reason why not to have liquid clustering for Iceberg.

Speaker 2 [40:33]

And if you had the chance today to magically migrate your data stack to Snowflake or Databricks, which one would you pick or both?

Speaker 1 [40:43]

Answer me, ask me afterwards.

Speaker 2 [40:46]

Not on stage. Okay, okay. So is there any reason why you chose specifically DynamoDB for cataloging?

Speaker 1 [40:54]

We didn't choose it, I just mentioned it, that you can use it.

Speaker 2 [40:54]

We didn't choose DynamoDB.

Speaker 1 [40:57]

We're not using DynamoDB for cataloging. If you require a locking mechanism, then you can use it. We don't use DynamoDB for locking or as a catalog.

Speaker 2 [41:06]

And, yeah, so I think another question goes in a similar direction. If you have different data in Delta and in Iceberg, and if yes, how do you tackle data governance? And if no, how do you keep indices in sync?

Speaker 1 [41:25]

So if you have both Delta Lake and Iceberg and how to keep them in sync or didn't quite get the question. Yeah

Speaker 2 [41:31]

Yeah, how to keep them in sync and how do you have data governance, because you need to account for both platforms.

Speaker 1 [41:38]

So ideally you have a single data catalog and this data catalog is independent of a cloud window or Like let's say Databricks or Snowflake think about Polaris even though it is a strong affiliation with Snowflake Lakekeeper is a great example a data catalog written in Rust and then with Lakekeeper as a data catalog you can manage both your iceberg and Your Delta Lake tables. There is no contradiction in doing so with both

Speaker 2 [42:08]

Why didn't you add the session boundaries into the data instead of having the boundaries as a file it could be identified by an extra column?

Speaker 1 [42:16]

Yeah, good question. We do have the session boundaries within the file itself, too, which is just the name of the log file.

Speaker 2 [42:16]

Yeah, good question.

Speaker 1 [42:22]

However, if you have it within the session file, then you have to take care yourself within the streaming context to remain history and to track history and to decide at which point is this session complete. So this is, once again, you start your manual custom framework in order to preserve these session boundaries. And that's something we didn't want to do. We just wanted an out-of-the-box solution, but this failed. And if you do this with Spark streaming, maintaining your own Spark context, your own history, then you get very close to what we had before and we wanted to ditch this.

Speaker 2 [42:56]

And we still have time for last two questions. So why do you create the session data on customer appliances instead of streaming the data directly via AWS, Firehose, Azure Event Hubs, or other?

Speaker 1 [43:09]

Yet another great question. Actually, we do both. We already stream the data, but we only stream high level KPI data, telemetry data. But with those log sessions, they contain very low level information. You can, like, pin a certain module in a given printing unit very down to, like, on some board where this message originated from. And if we would stream all the data that we have in our log sessions, this would be just too expensive. And there's also a technological point that this may not even be possible with our architecture to stream everything. So it's a trait of batch data for things that are not time-critical, but we stream data for time-critical applications.

Speaker 2 [43:54]

And then as a last question, what are your thoughts on S3 tables?

Speaker 1 [43:58]

Well, I think it's a good decision. They went for three tables by AWS in the beginning. I was a bit disappointed because at three tables, they promised interoperability, but in the beginning, you could only use it with a custom Spark wrapper given by AWS or created by AWS, or you had to use the AWS Clue Catalog. So you're forced to use the Clue Catalog. You can't use Polaris or Lakekeeper to connect to the data. But they've changed this just recently, I think one or two weeks ago, that now the S3 tables also support natively the Iceberg REST API for data exchange. This is really great. I think this will be a game changer, even though also to add the functionality of the S3 tables in maintenance, like compaction is still very limited. It only supports bit packing. It does not support sorting or Z order. And it also does not support incremental compaction. So typically what you want to do, I just compact the data I got from yesterday, but not the entire table. What happens if you change the size, the target size? Does it compact the entire table again? This would be hilarious. But you only want to compact the data that arrived like yesterday and not more. And this is also not possible at the moment. So it's still limited, but I think the direction is good. I think cost-wise, I can't really judge, but from what I've read, it's fairly expensive in comparison if you do it on yourself or by yourself.

Speaker 2 [45:20]

So thank you for the talk and for the answers France. There's still a few questions asked Maybe you can come to the front later or you can answer them on this code if you want or join

Speaker 1 [45:29]

Or join our booth, yeah.

Speaker 2 [45:29]

So, yeah, join him in his booth. So thank you for the presentation, and yeah, thank you. Thank you.

Franz Wöllert

About — in the speaker's own words

Hi my name is Franz and I’m an open source and python enthuisiast:

  • father of 3 girls
  • major in psychology
  • chess hobbiyst
  • competitive ultimate frisbee player
  • likes cooking and baking sourdough bread
Social card for talk: Open Table Formats in the Wild: From Parquet to Delta Lake and Back