Transformări de date în Databricks

Concepte Databricks

Kevin Barlow

Data Practitioner

SQL pentru ingineria datelor

SQL

  • Familiar pentru administratorii de baze de date (DBA)
  • Ideal pentru manipulări standard
  • Execută UDF-uri predefinite
-- Creating a new table in SQL

CREATE TABLE table_name
USING delta
AS (
  SELECT *
  FROM source_table
  WHERE date >= '2023-01-01'
)
Concepte Databricks

Alte limbaje pentru ingineria datelor

Python, R, Scala

  • Familiar inginerilor software
  • Transformări standard și complexe
  • Utilizarea și definirea funcțiilor personalizate
#Creating a new table in Pyspark

spark
  .read
  .table('source_table')
  .filter(col('date') >= '2023-01-01')
  .write
  .saveAsTable('table_name')
Concepte Databricks

Transformări comune

Manipularea schemei

  • Adăugarea și eliminarea coloanelor
  • Redefinirea coloanelor

#Pyspark

df
  .withColumn(col('newCol'), ...)
  .drop(col('oldCol'))

Filtrare

  • Reducerea DataFrame la un subset de date
  • Aplicarea mai multor criterii

#Pyspark

df
  .filter(col('date') >= target_date)
  .filter(col('id') IS NOT NULL)
Concepte Databricks

Transformări comune (continuare)

Date imbricate

  • Date de tip Array sau Struct
  • Extindere sau contracție
df
  .explode(col('arrayCol')) #wide to long
  .flatten(col('items')) #long to wide

Agregare

  • Gruparea datelor pe coloane
  • Calcularea rezumatelor de date
df
  .groupBy(col('region'))
  .agg(sum(col('sales')))
Concepte Databricks

Auto Loader

Auto Loader procesează fișierele noi pe măsură ce ajung în data lake.

  • Procesare incrementală
  • Procesare eficientă
  • Automatizat
spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .load(file_path)

Diagramă Auto Loader

1 https://www.databricks.com/blog/2020/02/24/introducing-databricks-ingest-easy-data-ingestion-into-delta-lake.html
Concepte Databricks

Structured Streaming

Pipeline de streaming

spark.readStream
    .format("kafka")
    .option("subscribe", "<topic>")
    .load()
    .join(table_df, 
      on="<id>", how="left")
    .writeStream
    .format("kafka")
    .option("topic", "<topic>")
    .start()
Concepte Databricks

Să exersăm!

Concepte Databricks

Preparing Video For Download...