Spark schema evolution in practice

Dear readers, good day!

In this article, the leading consultant of the Big Data Solutions business unit at "Neoflex" provides a detailed description of options for building variable structure data marts using Apache Spark.

In data analysis projects, the task of building data marts based on weakly structured data frequently arises.

Typically, this consists of logs or responses from various systems saved in JSON or XML format. The data is unloaded into Hadoop, and then a data mart needs to be constructed from it. We can organize access to the created data mart, for example, through Impala.

In this case, the schema of the target data mart is initially unknown. Moreover, the schema cannot be predefined because it depends on the data, and we are dealing with these very weakly structured data.

For instance, today the following response is logged:

{source: "app1", error_code: ""}

and tomorrow the same system sends the following response:

{source: "app1", error_code: "error", description: "Network error"}

As a result, another field — description — should be added to the data mart, and it is uncertain whether it will arrive or not.

The task of creating a data mart from such data is quite standard, and Spark offers several tools for this. For parsing the source data, support is available for both JSON and XML, and for unknown schemas, schemaEvolution support is provided.

At first glance, the solution appears straightforward. One needs to take a folder with JSON and read it into a dataframe. Spark will create a schema, converting nested data into structures. Then everything needs to be saved in parquet, which is also supported by Impala, by registering the data mart in the Hive metastore.

It seems simple enough.

However, the short examples in the documentation do not clarify what to do with a number of practical problems.

The documentation describes an approach not for creating a data mart but for reading JSON or XML into a dataframe.

Specifically, it simply explains how to read and parse JSON:

df = spark.read.json(path...)

This is sufficient to make the data accessible to Spark.

In practice, the scenario is much more complicated than simply reading JSON files from a folder and creating a dataframe. The situation is as follows: there is already a specific data mart, new data arrives every day, and it needs to be added to the data mart, bearing in mind that the schema may differ.

The usual schema for building a data mart is as follows:

Step 1. Data is loaded into Hadoop with subsequent daily incremental loads, resulting in a new partition. This creates a folder partitioned by day with the raw data.

Step 2. During the initial load, this folder is read and parsed using Spark. The resulting dataframe is saved in a format suitable for analysis, such as parquet, which can then be imported into Impala. Thus, a target view is created with all the data accumulated up to that point.

Step 3. A load is created that will update the view daily.
There arises the question of incremental loading, the need for partitioning the view, and maintaining an overall schema for the view.

For example, let's say the first step of the storage construction has been implemented, and the export of JSON files to the folder has been configured.

Creating a dataframe from them to save as a view is not an issue. This is the very first step, which can easily be found in the Spark documentation:

df = spark.read.option("mergeSchema", True).json(".../*") 
df.printSchema()

root 
|-- a: long (nullable = true) 
|-- b: string (nullable = true) 
|-- c: struct (nullable = true) |    
|-- d: long (nullable = true)

Everything seems fine.

We read and parsed the JSON, then saved the dataframe as parquet, registering it in Hive in any convenient way:

df.write.format(“parquet”).option('path','').saveAsTable('')

We obtain the view.

However, on the next day, new data from the source was added. We have a folder with JSON, and a view created based on that folder. After loading the next batch of data from the source, the view is missing data for one day.

A logical solution would be to partition the view by days, allowing a new partition to be added each subsequent day. This mechanism is also well-known; Spark allows writing partitions separately.

First, we perform the initial load, saving the data as described above, adding only partitioning. This action is called initializing the view and is done only once:

df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable)

The next day, we only load the new partition:

df.coalesce(1).write.mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable +"\/date_load=" + date_load + "\/" )

All that remains is to re-register in Hive to update the schema.
However, this is where problems arise.

The first problem. Sooner or later, the resulting parquet cannot be read. This is related to the different ways parquet and JSON handle empty fields.

Let's consider a typical situation. For example, yesterday, the JSON looked like this:

Day 1: {"a": {"b": 1}},

but today the same JSON looks like this:

Day 2: {"a": null}

Suppose we have two different partitions, each with one row.
When we read the raw data as a whole, Spark will be able to determine the type and understand that "a" is a field of type "struct", with a nested field "b" of type INT. However, if each partition was saved separately, it results in parquet files with incompatible partition schemas:

df1 (a: <struct>)
df2 (a: STRING NULLABLE)

This situation is well-known, so an option has been specifically added — when parsing the raw data, to remove empty fields:

df = spark.read.json("...", dropFieldIfAllNull=True)

In this case, the parquet will consist of partitions that can be read together.
However, those who have done this in practice will bitterly smile. Why? Because there will most likely arise two more situations. Or three. Or four. The first, which will almost certainly occur, is that numeric types will appear differently in different JSON files. For example, {intField: 1} and {intField: 1.1}. If such fields end up in the same partition, the schema merge will read everything correctly, casting to the most accurate type. But if they are in different partitions, one will have intField: int, while the other will have intField: double.

To handle this situation, there is the following flag:

df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)

Now we have a folder containing partitions that can be read into a single dataframe and a valid parquet of the entire showcase. Right? No.

We must remember that we registered the table in Hive. Hive is case-insensitive regarding field names, while parquet is case-sensitive. Therefore, partitions with schemas: field1: int, and Field1: int are the same for Hive but not for Spark. We need to ensure that field names are converted to lowercase.

After that, everything seems fine.

However, it’s not so simple. A second, also well-known problem arises. Since each new partition is saved separately, there will be auxiliary Spark files in the partition folder, such as the success flag _SUCCESS. This will cause an error when trying to read parquet. To avoid this, you need to configure the settings to prevent Spark from writing auxiliary files into the folder:

hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")

It seems that every day a new parquet partition is added to the target showcase folder, containing parsed data for the day. We have taken precautions to avoid partitions with data type conflicts.

However, we now face a third problem. The overall schema is unknown, and moreover, the Hive table has an incorrect schema, as each new partition has likely distorted the schema.

We need to re-register the table. This can be done simply: re-read the parquet showcases, take the schema, and create DDL based on it to re-register the folder in Hive as an external table, updating the target showcase schema.

We encounter a fourth problem. When we registered the table for the first time, we relied on Spark. Now we are doing this ourselves and need to remember that parquet fields may start with characters that are not allowed in Hive. For instance, Spark discards rows that it cannot parse into the field "corrupt_record". Such a field cannot be registered in Hive without escaping.

Knowing this, we get the schema:

f_def = ""
for f in pf.dtypes:
  if f[0] != "date_load":
    f_def = f_def + "," + f[0].replace("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") 
table_define = "CREATE EXTERNAL TABLE jsonevolvtable (" + f_def[1:] + " ) "
table_define = table_define + "PARTITIONED BY (date_load string) STORED AS PARQUET LOCATION '\/user\/admin\/testJson\/testSchemaEvolution\/pq\/'"
hc.sql("drop table if exists jsonevolvtable")
hc.sql(table_define)

Code ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") creates a safe DDL, meaning instead of:

create table tname (_field1 string, 1field string)

With field names like "_field1, 1field", a safe DDL is created, where the field names are escaped: create table `tname` (`_field1` string, `1field` string).

The question arises: how to correctly obtain a dataframe with the complete schema (in pf code)? How do we get this pf? This is the fifth problem. Should we re-read the schema of all partitions from the folder with parquet files of the target showcase? This method is the safest, but heavy.

The schema already exists in Hive. A new schema can be obtained by combining the schema of the entire table with the new partition. This means we need to take the table's schema from Hive and combine it with the new partition's schema. This can be done by reading the test metadata from Hive, saving it to a temporary folder, and reading both partitions at once using Spark.

Essentially, we have everything we need: the original table schema in Hive and the new partition. We also have the data. All that's left is to obtain the new schema that combines the showcase schema and the new fields from the created partition:

from pyspark.sql import HiveContext
from pyspark.sql.functions import lit
hc = HiveContext(spark)
df = spark.read.json("...", dropFieldIfAllNull=True)
df.write.mode("overwrite").parquet(".../date_load=12-12-2019")
pe = hc.sql("select * from jsonevolvtable limit 1")
pe.write.mode("overwrite").parquet(".../fakePartiton/")
pf = spark.read.option("mergeSchema", True).parquet(".../date_load=12-12-2019/*", ".../fakePartiton/*")

Next, we create the DDL for registering the table, as in the previous fragment.
If the entire chain works correctly — that is, there was an initial load, and a properly created table exists in Hive — we obtain the updated table schema.

And lastly, the problem is that you cannot just add a partition to a Hive table, as it will become broken. It is necessary to instruct Hive to fix its partition structure:

from pyspark.sql import HiveContext
hc = HiveContext(spark) 
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)

A simple task of reading JSON and creating a showcase based on it turns into overcoming a series of implicit challenges, the solutions to which must be sought individually. Although these solutions are straightforward, they take a considerable amount of time to find.

To implement the showcase building, it was necessary to:

  • Add partitions to the showcase while removing service files
  • Deal with empty fields in the source data that Spark typed
  • Convert simple types to strings
  • Convert field names to lowercase
  • Separate data extraction from table registration in Hive (DDL creation)
  • Remember to escape field names that may be incompatible with Hive
  • Learn how to update the table registration in Hive

In conclusion, it should be noted that the solution for building showcases hides many pitfalls. Therefore, when difficulties arise in implementation, it is better to consult an experienced partner with successful expertise.

Thank you for reading this article, we hope the information proves to be useful.

Source: habr.com

Buy reliable website hosting with DDoS protection, VPS VDS servers 🔥 Buy reliable website hosting with DDoS protection, VPS VDS servers | ProHoster