21 ms·
Saving Millions by Dumping Java Serialization
- jnewhouse 9y agoAuthor here, let me know if you have any questions/want more details.
- user5994461 9y agoSo... what's quantcast?
- jnewhouse 9y agoWe're a big data advertise and measure company based in San Francisco. We run online display ad campaigns for marketers across realtime bidding exchanges (RTB), such as those run by Google and AppNexus. We also provide a publisher product to give site owners insights into their audience. Stack Overflow's profile is at https://www.quantcast.com/stackoverflow.com https://www.quantcast.com/stackoverflow.com.
- deleted 9y ago[deleted]
- whack 9y agoThanks for writing about your experiences. Why not use an existing serialization framework, such as Protobuf, instead of building something in-house?
- jnewhouse 9y agoI don't think protobuf was around for public use when we came up with this format, which began around 2005. We use Protobuf internally, and some of our columns are actually byte[]'s containing protobuf data. We now support Parquet and are doing more work with other big data tools, but we've had a hard time matching the performance of our custom stuff.
- revscat 9y ago1) Can you provide any more details about how Rowfiles are structured and/or implemented? Specifically, how does it handle nested objects? Does it support `transient`? Do `writeObject` and/or `readObject` come into play? 2) Do you feel this is a generic enough solution that you would consider submitting it as a JSR?
- jnewhouse 9y agoIt natively supports a limited set of Columns. Basically boxed primitives, java.util.Date, joda.time.DateTime, and arrays and double arrays of both boxed and unboxed versions of the preceding. The list of Columns being used is used to read and write to a byte buffer. The byte buffer is almost entirely the field's data, with one or two bytes describing how the subsequent field is encoded. Nested objects aren't handled out of the box, but there is the capability to define a UserRowField that allows for serialization/deserialization to bytes of any Serializable class. This gets used for our SQL map-reduce function a lot. The downside is that you need to have the UserRowField implementation in your classpath in order to read the Row, which is not generally the case.
- user123 9y agoTLDR: we had shitty code, optimized it, now it runs well. No code examples, nothing.
- jnewhouse 9y agoIf you want more details, we were packing a Row class into a base64 encoded string using an ObjectOutputStream. This is a fine thing for small scale serialization but sucks at scale, because of the reasons mentioned in the post. Sorry we don't have code examples, but it's unclear how useful it'd be given that no one else uses our file format. If you want a bit more detail on how the format works. Each metadata contains a list of typed columns to define the schema of a given part. Our map-reduce framework has a bunch of internal logic that tries to justify the written Row class with the one the Mapper class is asking for. This allows us to do things like ingest different versions of a row with in the context of a single job. I think questions of serialization at the scale are generally interesting, although ymmv. I know of one company using Avro, which doesn't let you cleanly update or track schema. They've ended up storing every schema in an HBase table and reserving the first 8 bytes to do a lookup into this table to know the row's schema.
- cakoose 9y agoWhat do you mean when you say Avro doesn't let you "cleanly update or track schema"? From what I've read about Avro 1. It can transform data between two compatible schemas. 2. It can serialize/load schemas off the wire, so you can send the schema in the header. If schema serialization causes too much overhead, you can set things up so you only send the schema version identifier, as long as the receiver can use that to get access to the full schema.
- jnewhouse 9y agoI think what I'd heard about was likely a poorly implemented use of Avro. I haven't actually worked with it.
- barrkel 9y agoAvro can store the schema inline or out of line; with inline schemas, it's at the start of the file (embedded JSON), and it describes the schema for all the rows in that file. If you're working with Hive, the schema you put in the Hive metastore is cross-checked with each Avro file read; if any given Avro file doesn't contain a particular column, it just turns up as null for that subset of rows. Spark and Impala work similarly. I agree serialization at scale is interesting. My particular interest right at this moment is in efficiently doing incremental updates of HDFS files (Parquet & Avro) from observing changes in MySQL tables - not completely trivial because some ETL with joins and unions is required to get data in the right shape.
- Alupis 9y ago> Secondly, Java serialization produces very bulky outputs. Each serialization contains all of the data required to deserialize. When you’re writing billions of records at a time, recording the schema in every record massively increases your data size. Sounds to me like you shouldn't be storing objects in your database. Why not just write the data into tables, and then create new POJO's when necessary, using the selected data?
- jnewhouse 9y agoA standard database table isn't large enough to handle our large datasets. For example, the Hercules dataset was over 2 petabytes and even after optimization is almost 1 petabyte. Big data systems like Spark, Impala, Presto, etc. are designed to make the data look like a table, even though it is spread out into many files in a distributed filesystem. This is what we do. It's pretty common to reimplement some database features onto these big data file formats. In our case we have very fast indexes that let us quickly fetch data, similar to an index in a postgresql table.
- Alupis 9y agoWell, you understand your system and requirements better than I, obviously, but... A standard database table isn't large enough to handle our large datasets ... isn't much of an answer as-to why you're storing objects in your database. As you already mentioned in your post, serialized objects are big - they contain all of their data, plus everything necessary to deserialize the object into something usable. I imagine your objects have the standard amount of strings, characters, numbers, booleans, etc... why not just store those in the database and select them back out when needed? Less data in the database, and faster retrieval time since you skip serialization in both steps (storage and retrieval). Even if you have nested objects within nested objects, you can write-out a "flat" version of the data to a couple of joined tables surely. On the other hand, serializing the object is probably more "simple" to implement and use... but then you get the classical tradeoff of performance vs. convenience.
- barrkel 9y ago
- hrshtr 9y agoWas using Thrift or Protobuf an option?
- quest88 9y agoI'd like to know this too. As a passerby, those seem to have solved serialization, so I'm curious why you need rowfiles instead of e.g. protobuf.
- hrshtr 9y agoOne reason on top of my head: Using such communication protocol would require changes to the other services consuming it.
- rst 9y agoSo did switching to their homebrew serialization format -- in fact, most of the article is about how they managed the changes (which touched codebases at multiple sites in a fairly large organization).
- jnewhouse 9y agoThose switches all occurred at the pipeline level, leaving the map-reduce platform untouched. Switching our base logs to something like Parquet, Thrift or Protobuf would be a much larger project. We do support writing and reading Parquet to allow us to interface with other big data systems.
- pokemon-trainer 9y agoWhat is Thrift? Is it a service?
- givemefive 9y agoThrift is a software library not a BaaS...
- 9y ago
- MS_Buys_Upvotes 9y agoCan someone explain to an amateur why serialization is faster than say passing raw JSON? It seems like parsing JSON would be faster than the serialize -> deserialize process but with the popularity of things like Protobuff it's clear that JSON is slower.
- jerf 9y agoSerialization is the process of writing arbitrary data out into a blob of some sort (binary, text, whatever) that can be read in later and processed back into the original data, possibly not by the same system. This should be considered to include even the degenerate case of just writing the content of an expanse of RAM out, as that still raises issues related to serialization. "JSON Serialization" and "Java Serialization" are two different things that can accomplish that goal. It sounds to me from your question that you think they have some fundamental difference, because your second paragraph implies you believe there is some sort of fundamental difference between Java serialization and JSON serialization, but there isn't. There is a whole host of non-fundamental differences that you always have to consider with a serialization format (speed, what can be represented, circular data structure handling, whether untrusted data can be used), but there's not a fundamental difference.
- felixgallo 9y agoJSON must also be serialized or deserialized. Parsing it is slow and hard and not cache friendly. Protobuf has the benefit of being extremely compact and, in some important languages, fast and friendly to serialize and deserialize.
- MS_Buys_Upvotes 9y agoThanks! You and the others are right: I didn't know JSON was serialized. I can see why something in binary would be faster than structured text (think assembler vs Python). Thanks again.
- coldtea 9y ago>I didn't know JSON was serialized. Think of it like this: anytime you get stuff from the memory of your program (arrays, lists, strings, etc) and export it in a textual or binary format that can be exchanged between programs, moved over the network, saved to a file, etc, that's serialization.
- Cieplak 9y agoSome interesting benchmarks of various Java serialization libraries: https://github.com/eishay/jvm-serializers/wiki https://github.com/eishay/jvm-serializers/wiki
- jankotek 9y agoPerhaps I could share my project which is trying to 'fix' java serialization? It was originally part of database engine, but was extracted into separate project. It solves things like cyclic reference, non-recursive graph traversal and incremental serialization of large object graphs. https://github.com/jankotek/elsa/ https://github.com/jankotek/elsa/
- bluecarbuncle 9y agoPortable Object Format from 10 years ago? https://docs.oracle.com/cd/E24290_01/coh.371/e22837/api_pof.htm#COHDG1367 https://docs.oracle.com/cd/E24290_01/coh.371/e22837/api_pof....
- cntlzw 9y agoI am by no means an expert, but I always wonder why people don't adopt ASN.1 for serialization? I know it is not pretty but writing machine readable stuff never is.