Estimados lectores, ¡buen día!
En este artículo, el consultor líder de la línea de negocios Big Data Solutions de la compañía «Neoflex» describe detalladamente las opciones para construir vitrinas de estructura variable utilizando Apache Spark.
En el marco de un proyecto de análisis de datos, a menudo surge la tarea de construir vitrinas a partir de datos poco estructurados.
Generalmente, estos son registros o respuestas de diversos sistemas, almacenados en formato JSON o XML. Los datos se descargan en Hadoop, y luego hay que construir una vitrina a partir de ellos. Podemos organizar el acceso a la vitrina creada, por ejemplo, a través de Impala.
En este caso, el esquema de la vitrina objetivo es previamente desconocido. Además, el esquema no se puede elaborar de antemano, ya que depende de los datos y estamos tratando precisamente con esos datos poco estructurados.
Por ejemplo, hoy se registra la siguiente respuesta:
{source: "app1", error_code: ""}y mañana de este mismo sistema llega esta respuesta:
{source: "app1", error_code: "error", description: "Error de red"}Como resultado, se debe agregar un campo más a la vitrina: description, y si llegará o no, nadie lo sabe.
La tarea de crear una vitrina con tales datos es bastante estándar, y Spark tiene una serie de herramientas para ello. Para analizar los datos de origen, hay soporte tanto para JSON como para XML, y se prevé soporte para schemaEvolution para esquemas desconocidos.
A primera vista, la solución parece sencilla. Hay que tomar la carpeta con JSON y leerla en un dataframe. Spark creará el esquema, transformando los datos anidados en estructuras. Luego, todo debe guardarse en parquet, que también es compatible con Impala, registrando la vitrina en el Hive metastore.
Parece que todo es simple.
Sin embargo, de los breves ejemplos en la documentación no queda claro cómo abordar varios problemas en la práctica.
La documentación describe un enfoque no para crear una vitrina, sino para leer JSON o XML en un dataframe.
En particular, simplemente se indica cómo leer y parsear JSON:
df = spark.read.json(path...)Esto es suficiente para hacer que los datos sean accesibles para Spark.
En la práctica, el escenario es mucho más complicado que simplemente leer archivos JSON de una carpeta y crear un dataframe. La situación es la siguiente: ya existe una vitrina determinada, cada día llegan nuevos datos, hay que agregarlos a la vitrina, sin olvidar que el esquema puede diferir.
El esquema habitual para la construcción de vitrinas es el siguiente:
Paso 1. Los datos se cargan en Hadoop con una carga diaria posterior y se almacenan en una nueva partición. Se crea una carpeta particionada por días con datos originales.
Paso 2. Durante la carga inicial, esta carpeta se lee y se analiza utilizando herramientas de Spark. El dataframe resultante se guarda en un formato accesible para el análisis, como parquet, que luego se puede importar en Impala. Así se crea una vitrina de destino con todos los datos acumulados hasta ese momento.
Paso 3. Se crea una carga que actualizará la vitrina todos los días.
Surge la cuestión de la carga incremental, la necesidad de particionar la vitrina y la pregunta sobre el soporte del esquema general de la vitrina.
Veamos un ejemplo. Supongamos que se ha implementado el primer paso en la construcción del almacenamiento y se ha configurado la exportación de archivos JSON en la carpeta.
Crear un dataframe a partir de ellos para luego guardarlo como vitrina no supone ningún problema. Este es el primer paso que se puede encontrar fácilmente en la documentación de Spark:
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)Todo parece estar bien.
Hemos leído y analizado el JSON, luego guardamos el dataframe como parquet, registrándolo en Hive de cualquier forma conveniente:
df.write.format(“parquet”).option('path','').saveAsTable('') Obtenemos la vitrina.
Pero, al día siguiente, se agregaron nuevos datos de la fuente. Tenemos una carpeta con JSON y una vitrina creada a partir de esta carpeta. Después de cargar el siguiente lote de datos de la fuente, faltan datos por un día en la vitrina.
La solución lógica será particionar la vitrina por días, lo que permitirá agregar una nueva partición cada día siguiente. Este mecanismo también es bien conocido, Spark permite guardar particiones por separado.
Primero hacemos una carga inicial, guardando los datos como se describió anteriormente, añadiendo solo la partición. Esta acción se llama inicialización de la vitrina y se realiza solo una vez:
df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable)
Al día siguiente, solo cargamos la nueva partición:
df.coalesce(1).write.mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable +"\/date_load=" + date_load + "\/" )
Solo queda volver a registrar en Hive para actualizar el esquema.
Sin embargo, aquí es donde surgen problemas.
El primer problema. Tarde o temprano, el parquet generado no será legible. Esto se debe a cómo parquet y JSON manejan los campos vacíos de manera diferente.
Consideremos una situación típica. Por ejemplo, ayer llegó un JSON:
Día 1: {"a": {"b": 1}},
y hoy el mismo JSON se ve así:
Día 2: {"a": null}
Supongamos que tenemos dos particiones diferentes, cada una con una fila.
Cuando leemos los datos originales en su totalidad, Spark podrá determinar el tipo y entenderá que «a» es un campo del tipo «estructura», con un campo anidado «b» de tipo INT. Pero, si cada partición se guardó por separado, se obtiene un parquet con esquemas de partición incompatibles:
df1 (a: <struct>)
df2 (a: STRING NULLABLE)
Esta situación es bien conocida, por lo que se ha agregado específicamente la opción de eliminar los campos vacíos al analizar los datos originales:
df = spark.read.json("...", dropFieldIfAllNull=True)
En este caso, el parquet consistirá en particiones que se pueden leer juntas.
Aunque aquellos que lo han hecho en la práctica aquí se reirán amargamente. ¿Por qué? Porque probablemente surgirán otras dos situaciones. O tres. O cuatro. La primera, que casi seguramente ocurrirá, es que los tipos numéricos aparecerán de diferentes maneras en diferentes archivos JSON. Por ejemplo, {intField: 1} y {intField: 1.1}. Si tales campos caen en una misma partición, la fusión de esquemas los leerá correctamente, convirtiéndolos al tipo más preciso. Pero si están en diferentes, entonces en uno será intField: int, y en el otro intField: double.
Para manejar esta situación hay el siguiente flag:
df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)
Ahora tenemos una carpeta donde están las particiones que se pueden leer en un solo dataframe y un parquet válido para toda la vitrina. ¿Verdad? No.
Hay que recordar que registramos la tabla en Hive. Hive no distingue entre mayúsculas y minúsculas en los nombres de los campos, mientras que parquet sí lo hace. Así, las particiones con esquemas: field1: int, y Field1: int son iguales para Hive, pero no para Spark. No hay que olvidar convertir los nombres de los campos a minúsculas.
Después de esto, parece que todo está bien.
Sin embargo, no todo es tan simple. Surge un segundo problema, también bien conocido. Dado que cada nueva partición se guarda por separado, en la carpeta de la partición habrá archivos de servicio de Spark, como el flag de éxito de operación _SUCCESS. Esto conducirá a un error al intentar parquet. Para evitar esto, hay que configurar la configuración para impedir que Spark escriba archivos de servicio en la carpeta:
hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")
Parece que ahora cada día se agrega una nueva partición parquet a la carpeta del escaparate de destino, donde se encuentran los datos analizados del día. Nos hemos preocupado por anticipado de que no haya particiones con conflictos de tipos de datos.
Pero, enfrente de nosotros está el tercer problema. Ahora, el esquema general no es conocido, además, en Hive la tabla tiene un esquema incorrecto, ya que cada nueva partición probablemente ha distorsionado el esquema.
Es necesario volver a registrar la tabla. Esto se puede hacer de manera sencilla: volver a leer el escaparate parquet, tomar el esquema y crear a partir de él un DDL, con el cual volver a registrar la carpeta en Hive como una tabla externa, actualizando así el esquema del escaparate de destino.
Surge un cuarto problema. Cuando registramos la tabla por primera vez, nos basamos en Spark. Ahora lo hacemos nosotros mismos y debemos recordar que los campos parquet pueden comenzar con caracteres no permitidos para Hive. Por ejemplo, Spark descarta las filas que no pudo analizar en el campo «corrupt_record». Un campo así no podrá ser registrado en Hive sin realizar un escape.
Conociendo esto, obtenemos el esquema:
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)
Código ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace(«array<`», «array<») hace un DDL seguro, es decir, en lugar de:
create table tname (_field1 string, 1field string)
Con nombres de campos como «_field1, 1field», se hace un DDL seguro, donde los nombres de los campos están escapados: create table `tname` (`_field1` string, `1field` string).
Surge la pregunta: ¿cómo obtener correctamente un dataframe con el esquema completo (en el código pf)? ¿Cómo obtener este pf? Este es el quinto problema. ¿Leer de nuevo el esquema de todas las particiones de la carpeta con archivos parquet del escaparate de destino? Este método es el más seguro, pero también el más pesado.
El esquema ya está en Hive. Se puede obtener un nuevo esquema combinando el esquema de toda la tabla y la nueva partición. Esto significa que se debe tomar el esquema de la tabla desde Hive y combinarlo con el esquema de la nueva partición. Esto se puede hacer leyendo los metadatos de prueba desde Hive, guardándolos en una carpeta temporal y leyendo ambas particiones a la vez usando Spark.
En esencia, tenemos todo lo necesario: el esquema original de la tabla en Hive y la nueva partición. También tenemos los datos. Solo queda obtener un nuevo esquema que combine el esquema de la vitrina y los nuevos campos de la partición creada:
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/*")
A continuación, creamos el DDL para el registro de la tabla, como en el fragmento anterior.
Si toda la cadena funciona correctamente, es decir, se realizó la carga inicial y la tabla creada en Hive es correcta, entonces obtenemos el esquema actualizado de la tabla.
Y el último problema es que no se puede simplemente agregar una partición a la tabla Hive, ya que se romperá. Es necesario obligar a Hive a reparar su estructura de particiones:
from pyspark.sql import HiveContext
hc = HiveContext(spark)
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)
Una tarea simple de lectura de JSON y creación de una vitrina resulta en la superación de una serie de dificultades implícitas, soluciones para las cuales deben buscarse por separado. Y aunque estas soluciones son simples, se tarda mucho tiempo en encontrarlas.
Para implementar la construcción de la vitrina, fue necesario:
- Agregar particiones a la vitrina, deshaciéndose de los archivos temporales
- Resolver los campos vacíos en los datos de origen, que Spark tipificó
- Convertir tipos simples en cadenas
- Convertir los nombres de los campos a minúsculas
- Separar la exportación de datos del registro de la tabla en Hive (creación del DDL)
- No olvidar escapar los nombres de los campos que pueden ser incompatibles con Hive
- Aprender a actualizar el registro de la tabla en Hive
En conclusión, señalamos que la solución para construir vitrinas esconde muchos obstáculos. Por lo tanto, ante dificultades en la implementación, es mejor recurrir a un socio experimentado con una experiencia exitosa.
Gracias por leer este artículo, esperamos que la información sea útil.
Fuente: habr.com
Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster
