, Databricks ha vuelto a sacudir el mercado de datos. La empresa lanzó su edición gratuita de la plataforma Databricks con todas las funcionalidades incluidas. Es un gran recurso para aprender y probar, por decir lo menos.
Con eso en mente, creé un proyecto integral para ayudarte a aprender los fundamentos de los principales recursos dentro de Databricks.
Este proyecto demuestra un flujo de trabajo completo de extracción, transformación y carga (ETL) dentro de Databricks. Integra la API OpenWeatherMap para la recuperación de datos y el modelo OpenAI GPT-4o-mini para proporcionar sugerencias de vestimenta personalizadas basadas en el clima.
Aprendamos más al respecto.
El proyecto
El proyecto implementa una canalización de datos completa dentro de Databricks, siguiendo estos pasos.
Extracto: obtiene datos meteorológicos actuales de la ciudad de Nueva York a través de la API OpenWeatherMap[1]. Transformar: convierte marcas de tiempo UTC a la hora local de Nueva York y utiliza OpenAI.[2]GPT-4o-mini para generar sugerencias de aderezo personalizadas en función de la temperatura. Cargar: conserva los datos en Databricks Unity Catalog como archivos JSON sin formato y una tabla Delta estructurada (Capa plateada). Orquestación: el cuaderno con este código ETL se agrega a un trabajo y se programa para ejecutarse cada hora en Databricks. Análisis: la capa plateada alimenta un panel de Databricks que muestra información meteorológica relevante junto con las sugerencias del LLM.
Aquí está la arquitectura.
Excelente. Ahora que entendemos lo que debemos hacer, sigamos con el cómo de este tutorial.
Nota: si aún no tiene una cuenta en Databricks, vaya a la página de Databricks Free Edition[3], haga clic en Registrarse en la edición gratuita y siga las instrucciones que aparecen en pantalla para obtener acceso gratuito.
Extracto: Integración de API y Databricks
Como suelo decir, un proyecto de datos necesita datos para comenzar, ¿verdad? Entonces, nuestra tarea aquí es integrar la API OpenWeatherMap para ingerir datos directamente en un cuaderno PySpark dentro de Databricks. Esta tarea puede parecer complicada al principio, pero créanme, no lo es.
En la página inicial de Databricks, cree un nuevo cuaderno con el botón +Nuevo y luego seleccione Cuaderno.
Para la parte de Extracto, necesitaremos:
1. La clave API de API OpenWeatherMap.
Para conseguirlo, vaya a la página de registro de la API y complete su proceso de registro gratuito. Una vez que haya iniciado sesión en el panel, haga clic en la pestaña Clave API, donde podrá verla.
2. Importar paquetes
# Importaciones solicitudes de importación importar json
A continuación, crearemos una clase de Python para modularizar nuestro código y prepararlo también para producción.
Esta clase recibe la API_KEY que acabamos de crear, así como la ciudad y el país para la búsqueda del clima. Devuelve la respuesta en formato JSON. # Creando una clase para modularizar nuestra clase de código Weather: # Definir el constructor def __init__(self, API_KEY): self.API_KEY = API_KEY # Definir un método para recuperar datos meteorológicos def get_weather(self, city, country, unit='imperial'): self.city = city self.country = country self.units = unit # Realizar una solicitud GET a un punto final API que devuelve datos JSON url = f"https://api.openweathermap.org/data/2.5/weather?q={city},{country}&APPID={w.API_KEY}&units={units}" respuesta = request.get(url) # Utilice el método .json() para analizar el texto de respuesta y devolver si respuesta.status_code!= 200: elevar excepción(f"Error: {response.status_code} – {respuesta.texto}") devuelve respuesta.json()
Lindo. Ahora podemos ejecutar esta clase. Observe que usamos dbutils.widgets.get(). Este comando analiza los parámetros del trabajo programado, que veremos más adelante en este artículo. Es una buena práctica mantener los secretos a salvo.
# Obtener la clave API OpenWeatherMap API_KEY = dbutils.widgets.get('API_KEY') # Crear una instancia de la clase w = Weather(API_KEY=API_KEY) # Obtener los datos meteorológicos nyc = w.get_weather(city='New York', country='US') nyc
Aquí está la respuesta.
{'coord': {'lon': -74.006, 'lat': 40.7143}, 'weather': [{'id': 804, 'main': 'Nubes', 'description': 'nubes cubiertas', 'icon': '04d'}], 'base': 'estaciones', 'main': {'temp': 54.14, 'feels_like': 53,44, 'temp_min': 51,76, 'temp_max': 56,26, 'presión': 992, 'humedad': 89, 'sea_level': 992, 'grnd_level': 993}, 'visibilidad': 10000, 'viento': {'velocidad': 21,85, 'grados': 270, 'ráfaga': 37,98}, 'nubes': {'todos': 100}, 'dt': 1766161441, 'sys': {'tipo': 1, 'id': 4610, 'país': 'EE. UU.', 'amanecer': 1766146541, 'puesta del sol': 1766179850}, 'zona horaria': -18000, 'id': 5128581, 'nombre': 'Nueva York', 'cod': 200}
Con esa respuesta en la mano, podemos pasar a la parte de Transformación de nuestro proyecto, donde limpiaremos y transformaremos los datos.
Transformar: formatear los datos
En esta sección, veremos las tareas de limpieza y transformación realizadas sobre los datos sin procesar. Comenzaremos seleccionando los datos necesarios para nuestro panel. Se trata simplemente de obtener datos de un diccionario (o un JSON).
# Obteniendo información id = nyc['id'] marca de tiempo = nyc['dt'] clima = nyc['weather'][0]['main'] temp = nyc['main']['temp'] tmin = nyc['main']['temp_min'] tmax = nyc['main']['temp_max'] país = nyc['sys']['country'] ciudad = nyc['name'] amanecer = nyc['sys']['sunrise'] atardecer = nyc['sys']['sunset']
A continuación, transformemos las marcas de tiempo a la zona horaria de Nueva York, ya que viene con la hora de Greenwich.
# Transformar el amanecer y el atardecer en fecha y hora en la zona horaria de Nueva York desde fecha y hora importar fecha y hora, zona horaria desde zonainfo importar ZoneInfo importar hora # Marca de tiempo, amanecer y atardecer en la zona horaria de Nueva York target_timezone = ZoneInfo("America/New_York") dt_utc = datetime.fromtimestamp(sunrise, tz=timezone.utc) Sunrise_nyc = str(dt_utc.astimezone(target_timezone).time()) # get solo la hora del amanecer dt_utc = datetime.fromtimezone.time()) # obtener solo la hora del atardecer dt_utc = datetime.fromtimezone(timestamp, tz=timezone.utc) time_nyc = str(dt_utc.astimezone(target_timezone))
Finalmente, lo formateamos como un marco de datos Spark.
# Crear un marco de datos a partir de las variables df = spark.createDataFrame([[id, time_nyc, tiempo, temperatura, tmin, tmax, país, ciudad, amanecer_nyc, atardecer_nyc]], esquema=['id', 'timestamp','weather', 'temp', 'tmin', 'tmax', 'country', 'city', 'sunrise', 'sunset'])
El último paso en esta sección es agregar la sugerencia de un LLM. En este paso, seleccionaremos algunos de los datos obtenidos de la API y los pasaremos al modelo, pidiéndole que devuelva una sugerencia de cómo podría vestirse una persona para estar preparada para el clima.
Necesitará una clave API de OpenAI. Pase las condiciones climáticas, temperaturas máximas y mínimas (clima, tmax, tmin). Pídale al LLM que le devuelva una sugerencia sobre cómo vestirse según el clima. Agregue la sugerencia al marco de datos final. %pip install openai –quiet from openai import OpenAI import pyspark.sql.functions as F from pyspark.sql.functions import col # Obtener clave OpenAI OPENAI_API_KEY= dbutils.widgets.get('OPENAI_API_KEY') client = OpenAI( # Este es el valor predeterminado y se puede omitir api_key=OPENAI_API_KEY ) respuesta = client.responses.create( model="gpt-4o-mini", instrucciones="Eres un meteorólogo que da sugerencias sobre cómo vestirte según el clima. Responde en una oración.", input=f"El clima es {weather}, con temperatura máxima {tmax} y temperatura mínima {tmin}. ¿Cómo debo vestirme?" ) sugerencia = respuesta.output_text # Agrega la sugerencia a df df = df.withColumn('suggestion', F.lit(sugerencia)) display(df)
Fresco. Ya casi hemos terminado con el ETL. Ahora todo es cuestión de cargarlo. Esa es la siguiente sección.
Cargar: guardar los datos y crear la capa plateada
La última parte del ETL es cargar los datos. Lo cargaremos de dos formas diferentes.
Persistir los archivos sin formato en un volumen de catálogo de Unity. Guardar el marco de datos transformado directamente en la capa plateada, que es una tabla Delta lista para el consumo del Panel.
Creemos un catálogo que contendrá todos los datos meteorológicos que obtenemos de la API.
— Creando un catálogo CREAR CATÁLOGO SI NO EXISTE pipeline_weather COMENTARIO 'Este es el catálogo para el pipeline meteorológico';
A continuación, creamos un esquema para Lakehouse. Éste almacenará el volumen con los archivos JSON sin formato obtenidos.
— Creando un esquema CREAR ESQUEMA SI NO EXISTE pipeline_weather.lakehouse COMENTARIO 'Este es el esquema para el canal meteorológico';
Ahora, creamos el volumen para los archivos sin formato.
— Creemos un volumen CREAR VOLUMEN SI NO EXISTE pipeline_weather.lakehouse.raw_data COMENTARIO 'Este es el volumen de datos sin procesar para el canal meteorológico';
También creamos otro esquema para contener la tabla Delta de la capa plateada.
–Creando esquema para contener datos transformados CREAR ESQUEMA SI NO EXISTE pipeline_weather.silver COMENTARIO 'Este es el esquema para el canal meteorológico';
Una vez que tenemos todo configurado, así queda nuestro Catálogo.
Ahora, guardemos la respuesta JSON sin formato en nuestro volumen sin formato. Para mantener todo organizado y evitar la sobrescritura, adjuntaremos una marca de tiempo única a cada nombre de archivo.
Al agregar estos archivos al volumen en lugar de simplemente sobrescribirlos, estamos creando un "pista de auditoría" confiable. Esto actúa como una red de seguridad, lo que significa que si un proceso posterior falla o sufrimos una pérdida de datos más adelante, siempre podemos volver a la fuente y volver a procesar los datos originales cuando los necesitemos.
# Obtener marca de tiempo = datetime.now().strftime('%Y-%m-%d_%H-%M-%S') # Ruta para guardar json_path = f'/Volumes/pipeline_weather/lakehouse/raw_data/weather_{stamp}.json' # Guardar los datos en un archivo json df.write.mode('append').json(json_path)
Si bien mantenemos el JSON sin formato como nuestra "fuente de verdad", guardar los datos limpios en una tabla Delta en la capa Silver es donde ocurre la verdadera magia. Al utilizar .mode(“append”) y el formato Delta, nos aseguramos de que nuestros datos estén estructurados, aplicados por esquemas y listos para análisis de alta velocidad o herramientas de BI. Esta capa transforma las respuestas API desordenadas en una tabla confiable y consultable que crece con cada ejecución del proceso.
# Guarde los datos transformados en una tabla (esquema) ( df .write .format('delta') .mode("append") .saveAsTable('pipeline_weather.silver.weather') )
¡Hermoso! Con todo esto listo, veamos cómo luce nuestra mesa ahora.
Comencemos a automatizar este proceso ahora.
Orquestación: programación del portátil para que se ejecute automáticamente
Continuando con el proyecto, es hora de hacer que este oleoducto funcione por sí solo, con una mínima supervisión. Para eso, Databricks tiene la pestaña Trabajos y canalizaciones, donde es fácil programar trabajos para su ejecución.
Haga clic en la pestaña Trabajos y canalizaciones en el panel izquierdo. Busque el botón Crear y seleccione Trabajo. Haga clic en Notebook para agregarlo al trabajo. Configure como los datos a continuación. Agregue las claves API a los parámetros. Haga clic en Crear tarea. Haga clic en Ejecutar ahora para probar si funciona.
Una vez que haga clic en el botón Ejecutar ahora, debería comenzar a ejecutar el cuaderno y mostrar el mensaje Correcto.
Si el trabajo funciona bien, es hora de programarlo para que se ejecute automáticamente.
Haga clic en Agregar activador en el lado derecho de la pantalla, justo debajo de la sección Programaciones y activadores. Tipo de activador = Programado. Tipo de horario: seleccione Avanzado Seleccione Cada 1 hora en los menús desplegables. Guárdalo.
Excelente. ¡Nuestro Pipeline está en modo automático ahora! Cada hora, el sistema accederá a la API OpenWeatherMap y obtendrá información meteorológica actualizada para la ciudad de Nueva York y la guardará en nuestra Silver Layer Table.
Análisis: creación de un panel para decisiones basadas en datos
La última pieza de este rompecabezas es la creación del entregable de Analytics, que mostrará la información meteorológica y proporcionará al usuario información útil sobre cómo vestirse según el clima exterior.
Haga clic en la pestaña Paneles en el panel lateral izquierdo. Haga clic en el botón Crear panel. Se abrirá un lienzo en blanco para que podamos trabajar.
Ahora los paneles funcionan en función de los datos obtenidos de consultas SQL. Por lo tanto, antes de comenzar a agregar texto y gráficos al lienzo, primero debemos crear algunas métricas que serán las variables que alimentarán las tarjetas y los gráficos del tablero.
Entonces, haga clic en el botón +Crear desde SQL para iniciar una métrica. Dale un nombre. Por ejemplo, Ubicación, para recuperar el último nombre de la ciudad obtenida, debo usar la siguiente consulta.
— Obtener el último nombre de la ciudad obtenido SELECCIONAR ciudad DESDE pipeline_weather.silver.weather ORDENAR POR marca de tiempo DESC LIMIT 1
Y debemos crear una consulta SQL para cada métrica. Puedes verlos todos en el repositorio de GitHub [].
A continuación, hacemos clic en la pestaña Panel y comenzamos a arrastrar y soltar elementos en el lienzo.
Una vez que haces clic en el Texto, te permite insertar un cuadro en el lienzo y editar el texto. Cuando hace clic en el elemento gráfico, se inserta un marcador de posición para un gráfico y se abre el menú del lado derecho para seleccionar las variables y la configuración.
De acuerdo. Una vez agregados todos los elementos, el panel se verá así.
¡Qué lindo! Y así concluye nuestro proyecto.
Antes de ir
Puede replicar fácilmente este proyecto en aproximadamente una hora, según su experiencia con el ecosistema de Databricks. Si bien es una construcción rápida, incluye muchas habilidades de ingeniería básicas que podrás ejercitar:
Diseño arquitectónico: aprenderá cómo estructurar un entorno moderno de Lakehouse desde cero. Integración de datos perfecta: cerrará la brecha entre las API web externas y la plataforma Databricks para la ingesta de datos en tiempo real. Código limpio y modular: vamos más allá de los scripts simples mediante el uso de clases y funciones de Python para mantener la base de código organizada y mantenible. Automatización y orquestación: obtendrá experiencia práctica en la programación de trabajos para garantizar que su proyecto se ejecute de manera confiable en piloto automático. Ofrecer valor real: el objetivo no es sólo mover datos; es proporcionar valor. Al transformar las métricas meteorológicas sin procesar en sugerencias de vestimenta prácticas a través de IA, convertimos los "datos fríos" en un servicio útil para el usuario final.
Si te gustó este contenido, encuentra mis contactos y más sobre mí en mi sitio web.
https://gustavorsantos.me
Repositorio GitHub
Aquí está el repositorio de este proyecto.
https://github.com/gurezende/Databricks-Weather-Pipeline
Referencias
[1. API OpenWeatherMap] (https://openweathermap.org/)
[2. Plataforma abierta Ai] (https://platform.openai.com/)
[3. Edición gratuita de Databricks] (https://www.databricks.com/learn/free-edition)
[4. Repositorio de GitHub] (https://github.com/gurezende/Databricks-Weather-Pipeline)