In [1]:
pip install delta-spark==3.3.2 dotenv

Collecting pyspark<3.6.0,>=3.5.3 (from delta-spark==3.3.2)
  Using cached pyspark-3.5.7-py2.py3-none-any.whl
Installing collected packages: pyspark
  Attempting uninstall: pyspark
    Found existing installation: pyspark 3.5.0
    Can't uninstall 'pyspark'. No files were found to uninstall.
Successfully installed pyspark-3.5.7
Note: you may need to restart the kernel to use updated packages.


In [2]:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import *
from delta import configure_spark_with_delta_pip
from delta.tables import DeltaTable
from datetime import datetime

In [3]:
from dotenv import load_dotenv
import os
load_dotenv('/opt/workspace/.env')
MINIO_ENDPOINT = os.getenv("MINIO_ENDPOINT_DOCKER")
MINIO_ACCESS = os.getenv("MINIO_ROOT_USER")
MINIO_SECRET = os.getenv("MINIO_ROOT_PASSWORD")
                           
today = datetime.now().strftime("%Y/%m/%d")
BRONZE_PATH = f"s3a://bronze/posicao/{today}/"
SILVER_PATH = "s3a://silver/posicao/"
print(f"Buscando arquivos arquivos em: {BRONZE_PATH}")

Buscando arquivos arquivos em: s3a://bronze/posicao/2025/11/08/


In [None]:
builder = (
    SparkSession.builder.appName("BronzeToSilver_Delta")
    .config("spark.hadoop.fs.s3a.endpoint", MINIO_ENDPOINT)
    .config("spark.hadoop.fs.s3a.access.key", MINIO_ACCESS)
    .config("spark.hadoop.fs.s3a.secret.key", MINIO_SECRET)
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false")
    # Delta Lake
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
)
spark = configure_spark_with_delta_pip(builder).getOrCreate()

In [None]:
schema = StructType([
    StructField("hr", StringType(), True),
    StructField("l", ArrayType(
        StructType([
            StructField("c", StringType(), True),
            StructField("cl", IntegerType(), True),
            StructField("sl", IntegerType(), True),
            StructField("lt0", StringType(), True),
            StructField("lt1", StringType(), True),
            StructField("qv", IntegerType(), True),
            StructField("vs", ArrayType(
                StructType([
                    StructField("p", IntegerType(), True),
                    StructField("a", BooleanType(), True),
                    StructField("ta", StringType(), True),
                    StructField("py", DoubleType(), True),
                    StructField("px", DoubleType(), True),
                ])
            ), True)
        ])
    ), True)
])


In [None]:
# Pega o arquivo mais recente
files_df = spark.read.format("binaryFile").load(BRONZE_PATH)
latest_file_row = (
    files_df.orderBy(F.col("modificationTime").desc())
    .select("path")
    .limit(1)
    .collect()
)
latest_file = latest_file_row[0].path
print(latest_file)

In [None]:
df_raw = spark.read.option("mode", "PERMISSIVE").schema(schema).json(latest_file)

if df_raw.isEmpty() or df_raw.filter(F.col("hr").isNotNull() & F.col("l").isNotNull()).isEmpty():
    print("‚ö†Ô∏è Arquivo vazio.")
    spark.stop()
    exit(0)

In [None]:
# 1. EXPLODE o array "l" linha
df_linhas = df_raw.selectExpr("hr", "inline(l)")
df_linhas.show(2, truncate=False)
df_linhas.printSchema()

In [None]:
# 2. EXPLODE o array "vs" ve√≠culos
df_veiculos = df_linhas.selectExpr(
    "hr",
    "c as letreiro",
    "cl as codigo_linha",
    "sl as sentido",
    "lt0 as terminal_inicial",
    "lt1 as terminal_final",
    "qv",
    "inline(vs)"  # expande os ve√≠culos
)
df_veiculos.show(3, truncate=False)
df_veiculos.printSchema()

In [None]:
# 3. FILTRAR ONIBUS NAO REGULAR
df_filtrado = df_veiculos.filter(
    "codigo_linha IS NOT NULL AND NOT (codigo_linha < 1000 OR letreiro RLIKE 'GUIN|TEST|TST')"
)
print(f"Total de √¥nibus teste/guincho filtrados: {df_veiculos.count() - df_filtrado.count()}")

In [None]:
# 4. Selecionar e renomear colunas √∫teis
df_limpo = df_filtrado.select(
    "letreiro",
    "codigo_linha",
    "sentido",
    "terminal_inicial",
    "terminal_final",
    F.col("p").alias("codigo_veiculo"),
    F.col("a").alias("acessibilidade"),
    F.to_timestamp("ta").alias("ultima_atualizacao"),
    F.col("py").alias("latitude"),
    F.col("px").alias("longitude"),
    F.to_timestamp("hr").alias("hora_referencia"),
)
df_limpo.show(5, truncate=False)

In [None]:
# 5. Remover duplicatas de registros (se houver)
df_dedup = df_limpo.dropDuplicates(["codigo_veiculo", "hora_referencia"])
print(f"Quantidade de registros || Antes: {df_limpo.count()} | Depois: {df_dedup.count()} | # de duplicatas: {df_limpo.count() - df_dedup.count()}")

In [None]:
# 6. Corrigir imprecis√£o floating point lat/long e adicionar metadados
df_final = (
    df_dedup
    .withColumn("latitude", F.round("latitude", 6))
    .withColumn("longitude", F.round("longitude", 6))
    .withColumn("data_ref", F.to_date("ultima_atualizacao"))
    .withColumn("ingest_timestamp", F.current_timestamp())
)
df_final.show(3, truncate=False)
df_final.printSchema()

In [None]:
df_final.show(4, truncate=False)

In [None]:
if DeltaTable.isDeltaTable(spark, SILVER_PATH):
    silver_table = DeltaTable.forPath(spark, SILVER_PATH)
    (
        silver_table.alias("tgt")
        .merge(
            df_final.alias("src"),
            "tgt.codigo_veiculo = src.codigo_veiculo AND tgt.hora_referencia = src.hora_referencia"
        )
        .whenNotMatchedInsertAll()
        .execute()
    )
    print("üîÅ MERGE incremental conclu√≠do.")
else:
    (
        df_final.write.format("delta")
        .mode("overwrite")
        .partitionBy("data_ref")
        .save(SILVER_PATH)
    )
    print("üÜï Nova tabela Delta criada.")