Cómo convertir Parquet a Delta Lake en Databricks
Cómo convertir Parquet a Delta Lake en Databricks es una tarea común cuando se desea aprovechar las transacciones ACID, Time Travel y las optimizaciones de Delta. Esta guía práctica explica por qué y muestra paso a paso cómo transformar archivos Parquet en tablas Delta, con ejemplos en PySpark y comprobaciones sencillas.
Requisitos previos
- Workspace Databricks y un cluster ejecutando un Databricks Runtime con soporte para Delta.
- Permisos de lectura/escritura en el storage donde están los archivos Parquet (ej.: /mnt/...).
- Archivos Parquet de ejemplo y un notebook en Python (PySpark).
- Conocimientos básicos de Spark DataFrame y SQL.
Paso 1: Inspeccionar el archivo Parquet
Antes de convertir, confirme el esquema y posibles columnas de particionado. Conocer el esquema ayuda a decidir si necesita hacer alineación de esquema o renombrados.
# caminho do ficheiro Parquet
path_parquet = '/mnt/data/sample_parquet/'
# listar ficheiros e inspecionar schema
display(dbutils.fs.ls(path_parquet))
df = spark.read.parquet(path_parquet)
df.printSchema()
df.show(5, truncate=False)
Paso 2: Normalizar el esquema (opcional)
Si los archivos Parquet tienen columnas con nombres distintos o tipos inconsistentes, normalice antes de escribir. Ejemplo: uniformizar nombres y convertir tipos.
# exemplo simples de normalização
from pyspark.sql.functions import col
# renomear coluna e ajustar tipo se necessário
df = df.withColumnRenamed('OldName', 'new_name')
df = df.withColumn('event_date', col('event_date').cast('date'))
Paso 3: Escribir a Delta con particionado y evolución de esquema
Escriba el DataFrame a un directorio Delta. Use partitionBy para mejorar el rendimiento en lecturas, y opciones como mergeSchema para aceptar evoluciones de esquema cuando haga append.
path_delta = '/mnt/data/delta/sample_delta/'
# primeira gravação: criar a tabela Delta (overwrite se estiver a testar)
df.write.format('delta')
.mode('overwrite')
.option('overwriteSchema', 'true')
.partitionBy('year')
.save(path_delta)
# caso vá concatenar ficheiros com evolução de esquema, use:
# df_new.write.format('delta').mode('append').option('mergeSchema','true').save(path_delta)
Paso 4: Registrar la tabla Delta en el metastore (opcional pero recomendado)
Registrar la tabla facilita consultas SQL e integración con otras herramientas. Puede crear una tabla gestionada/externa que apunte a la ubicación Delta.
# criar esquema se necessário
spark.sql('CREATE DATABASE IF NOT EXISTS analytics')
# registar tabela externa que aponta para a pasta Delta
spark.sql("CREATE TABLE IF NOT EXISTS analytics.sample_delta USING DELTA LOCATION '/mnt/data/delta/sample_delta/'")
# agora pode consultar com SQL
spark.sql('SELECT COUNT(*) FROM analytics.sample_delta').show()
Paso 5: Optimizar y gestionar versiones (OPTIMIZE y VACUUM)
Después de cargar los datos, use OPTIMIZE para mejorar las lecturas y VACUUM para eliminar archivos antiguos (preste atención a la retención). Estos comandos son útiles en producción.
# OPTIMIZE (requer cluster com suporte Delta e permissões)
spark.sql('OPTIMIZE analytics.sample_delta')
# VACUUM com cuidado (padrão 7 dias); aqui um exemplo com 168 horas = 7 dias
spark.sql('VACUUM analytics.sample_delta RETAIN 168 HOURS')
Verificar el resultado
Confirme que la conversión fue correcta verificando la existencia del _delta_log, consultando recuentos y muestras, y listando la tabla en el metastore.
# verificar diretório Delta e _delta_log
display(dbutils.fs.ls(path_delta))
display(dbutils.fs.ls(path_delta + '/_delta_log'))
# ler a tabela Delta e mostrar linhas
spark.read.format('delta').load(path_delta).show(5)
# verificar tabela no metastore
spark.sql('SHOW TABLES IN analytics').show()
Errores comunes: permisos en el mount, columna de particionado ausente, conflicto de tipos cuando no se usa mergeSchema. Si ve exceptions sobre incompatibilidad de esquema, reevalúe la normalización del esquema o use mergeSchema con append.
Conclusión
Convertir Parquet a Delta Lake en Databricks permite aprovechar funcionalidades avanzadas como transacciones ACID y optimizaciones de lectura. Próximos pasos: automatizar el proceso con Workflows, integrar con Unity Catalog para gobernanza y probar OPTIMIZE con ZORDER. Consejo: antes de VACUUM, confirme la retención y haga un snapshot para evitar pérdida accidental de datos — ¿le apetece probar a convertir un conjunto de datos mayor para medir las ganancias de I/O?