Curso Prático de Engenharia de Dados Python
Curso Prático de Engenharia de Dados Python
Platform
Estrutura Pedagógica
Este curso segue uma progressão harmônica onde cada módulo representa 1.618x a complexidade do anterior,
respeitando a proporção áurea do aprendizado técnico.
Objetivos
Configurar ambiente local de desenvolvimento
Estabelecer conta GCP com free tier
Ferramentas
bash
# Instalações necessárias
python 3.11+
pip install virtualenv
gcloud SDK
docker desktop
git
VSCode ou PyCharm
PostgreSQL local
MongoDB Community Edition
python
# [Link]
import sys
import subprocess
import pkg_resources
def check_environment():
"""Verifica se o ambiente está configurado corretamente"""
checks = {
'Python': [Link],
'pip': [Link]('pip --version'),
'gcloud': [Link]('gcloud --version').split('\n')[0],
'docker': [Link]('docker --version'),
'postgres': [Link]('psql --version'),
'mongo': [Link]('mongod --version').split('\n')[0]
}
if __name__ == "__main__":
check_environment()
Conceitos Fundamentais
Estruturas de dados otimizadas
Manipulação eficiente com pandas
APIs REST
python
# weather_collector.py
import requests
import pandas as pd
from datetime import datetime
import sqlite3
from typing import Dict, List
class WeatherCollector:
"""Coleta dados meteorológicos e armazena localmente"""
def _init_database(self):
"""Inicializa banco SQLite com schema otimizado"""
conn = [Link](self.db_path)
[Link]("""
CREATE TABLE IF NOT EXISTS weather_data (
id INTEGER PRIMARY KEY AUTOINCREMENT,
city TEXT NOT NULL,
timestamp INTEGER NOT NULL,
temperature REAL,
humidity INTEGER,
pressure INTEGER,
description TEXT,
UNIQUE(city, timestamp)
)
""")
[Link]("CREATE INDEX IF NOT EXISTS idx_city_time ON weather_data(city, timestamp)")
[Link]()
[Link]()
if response.status_code == 200:
weather = [Link]()
[Link]({
'city': city,
'timestamp': int([Link]().timestamp()),
'temperature': weather['main']['temp'],
'humidity': weather['main']['humidity'],
'pressure': weather['main']['pressure'],
'description': weather['weather'][0]['description']
})
return [Link](data)
query = """
SELECT
AVG(temperature) as avg_temp,
MAX(temperature) as max_temp,
MIN(temperature) as min_temp,
AVG(humidity) as avg_humidity
FROM weather_data
WHERE city = ?
AND timestamp > strftime('%s', 'now', '-{} days')
""".format(days)
return df.to_dict('records')[0]
# Uso prático
if __name__ == "__main__":
collector = WeatherCollector()
# Coletar dados
cities = ["São Paulo", "Rio de Janeiro", "Brasília"]
df = collector.collect_weather(cities, "YOUR_API_KEY")
# Persistir
collector.save_to_database(df)
# Analisar
trends = collector.analyze_trends("São Paulo")
print(f"Tendências São Paulo: {trends}")
Conceitos
Design de schemas eficientes
Índices e performance
Window functions
sql
-- [Link]
CREATE SCHEMA IF NOT EXISTS sales_analytics;
python
# sales_analyzer.py
import psycopg2
import pandas as pd
from contextlib import contextmanager
from typing import Generator
class SalesAnalyzer:
"""Analisador de vendas com SQL otimizado"""
@contextmanager
def get_connection(self) -> Generator:
"""Context manager para conexões seguras"""
conn = [Link](self.connection_string)
try:
yield conn
finally:
[Link]()
cursor.copy_from(
buffer,
'temp_sales',
sep=',',
columns=['product_id', 'quantity', 'sale_date', 'customer_id']
)
[Link]()
Conceitos
Modelagem orientada a queries
Sharding e replicação
Índices em NoSQL
Patterns de agregação
python
# recommendation_engine.py
from pymongo import MongoClient, ASCENDING, TEXT
from datetime import datetime, timedelta
import numpy as np
from typing import List, Dict, Optional
from bson import ObjectId
class RecommendationEngine:
"""Motor de recomendação usando MongoDB"""
def _setup_collections(self):
"""Configura coleções com índices otimizados"""
# Coleção de usuários
if 'users' not in [Link].list_collection_names():
[Link].create_collection('users')
[Link].create_index([('email', ASCENDING)], unique=True)
[Link].create_index([('created_at', ASCENDING)])
# Coleção de interações
if 'interactions' not in [Link].list_collection_names():
[Link].create_collection('interactions')
[Link].create_index([
('user_id', ASCENDING),
('timestamp', ASCENDING)
])
[Link].create_index([
('product_id', ASCENDING),
('type', ASCENDING)
])
# Índice composto para queries de agregação
[Link].create_index([
('user_id', ASCENDING),
('product_id', ASCENDING),
('type', ASCENDING)
])
interaction = {
'user_id': ObjectId(user_id),
'product_id': ObjectId(product_id),
'type': interaction_type, # view, click, purchase, rating
'timestamp': [Link](),
'metadata': metadata or {}
}
pipeline = [
# Encontrar produtos que o usuário interagiu
{
'$match': {
'user_id': ObjectId(user_id)
}
},
# Agrupar por produto
{
'$group': {
'_id': '$product_id',
'user_score': {'$sum': '$score'}
}
},
# Encontrar outros usuários que interagiram com os mesmos produtos
{
'$lookup': {
'from': 'interactions',
'localField': '_id',
'foreignField': 'product_id',
'as': 'other_users'
}
},
# Expandir array de outros usuários
{
'$unwind': '$other_users'
},
# Filtrar o próprio usuário
{
'$match': {
'other_users.user_id': {'$ne': ObjectId(user_id)}
}
},
# Encontrar produtos que esses usuários também gostaram
{
'$lookup': {
'from': 'interactions',
'let': {'other_user': '$other_users.user_id'},
'pipeline': [
{
'$match': {
'$expr': {
'$and': [
{'$eq': ['$user_id', '$$other_user']},
{'$gte': ['$score', 3]}
]
}
}
}
],
'as': 'recommendations'
}
},
# Expandir recomendações
{
'$unwind': '$recommendations'
},
# Agrupar e calcular score de recomendação
{
'$group': {
'_id': '$recommendations.product_id',
'recommendation_score': {
'$sum': {
'$multiply': [
'$user_score',
'$[Link]',
'$other_users.score'
]
}
},
'support': {'$sum': 1} # Número de usuários que recomendam
}
},
# Filtrar produtos já vistos
{
'$lookup': {
'from': 'interactions',
'let': {'prod_id': '$_id'},
'pipeline': [
{
'$match': {
'$expr': {
'$and': [
{'$eq': ['$product_id', '$$prod_id']},
{'$eq': ['$user_id', ObjectId(user_id)]}
]
}
}
}
],
'as': 'already_seen'
}
},
{
'$match': {
'already_seen': {'$size': 0}
}
},
# Buscar detalhes do produto
{
'$lookup': {
'from': 'products',
'localField': '_id',
'foreignField': '_id',
'as': 'product'
}
},
{
'$unwind': '$product'
},
# Ordenar por score e limitar
{
'$sort': {
'recommendation_score': -1,
'support': -1
}
},
{
'$limit': limit
},
# Formatar saída
{
'$project': {
'product_id': '$_id',
'name': '$[Link]',
'category': '$[Link]',
'score': '$recommendation_score',
'recommended_by': '$support'
}
}
]
return list([Link](pipeline))
return list([Link](pipeline))
Conceitos
Arquitetura serverless
Monitoramento e logging
python
# cloud_functions/data_ingestion/[Link]
import functions_framework
from [Link] import storage, bigquery, firestore
import pandas as pd
import json
from datetime import datetime
from typing import Dict, Any
# Inicializar clientes
storage_client = [Link]()
bigquery_client = [Link]()
firestore_client = [Link]()
@functions_framework.cloud_event
def process_uploaded_file(cloud_event):
"""
Cloud Function triggerada por upload no Cloud Storage
Processa arquivo e carrega no BigQuery
"""
try:
# Baixar arquivo do Storage
bucket = storage_client.bucket(bucket_name)
blob = [Link](file_name)
# Adicionar metadados
df['processed_at'] = [Link]()
df['source_file'] = file_name
# Validação básica
df = validate_and_clean_data(df)
# Carregar no BigQuery
dataset_id = 'raw_data'
table_id = extract_table_name(file_name)
table_ref = f"{bigquery_client.project}.{dataset_id}.{table_id}"
job_config = [Link](
write_disposition=[Link].WRITE_APPEND,
schema_update_options=[
[Link].ALLOW_FIELD_ADDITION
]
)
job = bigquery_client.load_table_from_dataframe(
df, table_ref, job_config=job_config
)
[Link]() # Aguardar conclusão
# Registrar no Firestore
doc_ref = firestore_client.collection('processing_log').document()
doc_ref.set({
'file_name': file_name,
'records_processed': len(df),
'status': 'success',
'processed_at': [Link](),
'bigquery_table': table_ref
})
except Exception as e:
# Log de erro no Firestore
error_ref = firestore_client.collection('processing_errors').document()
error_ref.set({
'file_name': file_name,
'error': str(e),
'timestamp': [Link]()
})
raise e
return df
# [Link]
"""
functions-framework==3.*
google-cloud-storage==2.10.*
google-cloud-bigquery==3.11.*
google-cloud-firestore==2.11.*
pandas==2.0.*
pyarrow==12.0.*
"""
python
# app_engine/[Link]
from flask import Flask, request, jsonify
from [Link] import bigquery
from [Link] import firestore
import pandas as pd
from datetime import datetime, timedelta
import numpy as np
from functools import lru_cache
import json
app = Flask(__name__)
bigquery_client = [Link]()
firestore_client = [Link]()
cache_ref = firestore_client.collection('query_cache').document(query_hash)
cache_doc = cache_ref.get()
if cache_doc.exists:
cache_data = cache_doc.to_dict()
cache_time = cache_data.get('timestamp', [Link])
return None
@[Link]('/api/v1/analytics/sales', methods=['GET'])
def sales_analytics():
"""Endpoint para análise de vendas"""
# Parâmetros da query
start_date = [Link]('start_date',
([Link]() - timedelta(days=30)).strftime('%Y-%m-%d'))
end_date = [Link]('end_date',
[Link]().strftime('%Y-%m-%d'))
granularity = [Link]('granularity', 'daily') # daily, weekly, monthly
metrics = [Link]('metrics') or ['revenue', 'orders', 'avg_order_value']
# Query BigQuery
query = f"""
WITH sales_data AS (
SELECT
DATE(sale_timestamp) as sale_date,
order_id,
customer_id,
total_amount,
items_count
FROM `{bigquery_client.project}.[Link]`
WHERE DATE(sale_timestamp) BETWEEN @start_date AND @end_date
),
aggregated AS (
SELECT
{get_date_trunc(granularity)} as period,
COUNT(DISTINCT order_id) as orders,
COUNT(DISTINCT customer_id) as unique_customers,
SUM(total_amount) as revenue,
AVG(total_amount) as avg_order_value,
SUM(items_count) as total_items
FROM sales_data
GROUP BY period
)
SELECT
period,
{', '.join([f'{m} as {m}' for m in metrics if m in
['orders', 'unique_customers', 'revenue', 'avg_order_value', 'total_items']])}
FROM aggregated
ORDER BY period
"""
job_config = [Link](
query_parameters=[
[Link]("start_date", "DATE", start_date),
[Link]("end_date", "DATE", end_date)
]
)
response = {
'data': data,
'metadata': {
'start_date': start_date,
'end_date': end_date,
'granularity': granularity,
'metrics': metrics,
'generated_at': [Link]().isoformat()
}
}
# Salvar no cache
cache_ref = firestore_client.collection('query_cache').document(query_hash)
cache_ref.set({
'result': response,
'timestamp': [Link]()
})
return jsonify(response)
@[Link]('/api/v1/analytics/forecast', methods=['POST'])
def sales_forecast():
"""Previsão de vendas usando histórico"""
data = request.get_json()
product_id = [Link]('product_id')
days_ahead = [Link]('days_ahead', 7)
# Query histórico
query = """
SELECT
DATE(sale_timestamp) as date,
SUM(quantity) as daily_sales
FROM `{}.analytics.sales_items`
WHERE product_id = @product_id
AND DATE(sale_timestamp) >= DATE_SUB(CURRENT_DATE(), INTERVAL 90 DAY)
GROUP BY date
ORDER BY date
""".format(bigquery_client.project)
job_config = [Link](
query_parameters=[
[Link]("product_id", "STRING", product_id)
]
)
df = bigquery_client.query(query, job_config=job_config).to_dataframe()
# Calcular tendência
df['rolling_mean_7'] = df['daily_sales'].rolling(window=7).mean()
df['rolling_mean_30'] = df['daily_sales'].rolling(window=30).mean()
# Previsão ponderada
forecast = (recent_avg * 0.5 + medium_avg * 0.3 + long_avg * 0.2)
# Gerar previsões
future_dates = pd.date_range(
start=[Link]() + timedelta(days=1),
periods=days_ahead,
freq='D'
)
predictions = []
for date in future_dates:
# Adicionar sazonalidade semanal
day_of_week = [Link]
seasonality_factor = 1.0
[Link]({
'date': [Link]('%Y-%m-%d'),
'predicted_sales': round(forecast * seasonality_factor, 2),
'confidence_interval': {
'lower': round(forecast * seasonality_factor * 0.8, 2),
'upper': round(forecast * seasonality_factor * 1.2, 2)
}
})
return jsonify({
'product_id': product_id,
'forecast': predictions,
'model_metrics': {
'historical_average': round(df['daily_sales'].mean(), 2),
'recent_trend': 'increasing' if recent_avg > long_avg else 'decreasing',
'volatility': round(df['daily_sales'].std(), 2)
}
})
@[Link](Exception)
def handle_error(error):
"""Handler global de erros"""
# Log no Firestore
error_ref = firestore_client.collection('api_errors').document()
error_ref.set({
'endpoint': [Link],
'method': [Link],
'error': str(error),
'timestamp': [Link]()
})
return jsonify({
'error': 'Internal server error',
'message': str(error) if [Link] else 'An error occurred'
}), 500
if __name__ == '__main__':
[Link](host='[Link]', port=8080)
# [Link]
"""
runtime: python311
instance_class: F2
automatic_scaling:
target_cpu_utilization: 0.65
min_instances: 1
max_instances: 10
env_variables:
GAE_USE_SOCKETS_FOR_CLOUDSQL: 'true'
handlers:
- url: /api/.*
script: auto
secure: always
"""
python
# vertex_ai_pipeline/[Link]
from [Link] import aiplatform
from [Link] import bigquery
from [Link] import storage
import pandas as pd
import numpy as np
from sklearn.model_selection import train_test_split
from [Link] import StandardScaler
from [Link] import RandomForestRegressor
import joblib
from datetime import datetime
import json
class MLPipeline:
"""Pipeline de ML integrado com GCP"""
[Link](project=project_id, location=location)
self.bq_client = [Link](project=project_id)
self.storage_client = [Link](project=project_id)
query = f"""
WITH features AS (
SELECT
customer_id,
COUNT(DISTINCT order_id) as total_orders,
SUM(total_amount) as lifetime_value,
AVG(total_amount) as avg_order_value,
MAX(order_date) as last_order_date,
MIN(order_date) as first_order_date,
COUNT(DISTINCT product_category) as categories_purchased,
COUNT(DISTINCT DATE_TRUNC(order_date, MONTH)) as active_months
FROM `{self.project_id}.{dataset_id}.{table_id}`
GROUP BY customer_id
),
target AS (
SELECT
customer_id,
SUM(total_amount) as next_month_value
FROM `{self.project_id}.{dataset_id}.{table_id}`
WHERE order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 30 DAY)
GROUP BY customer_id
)
SELECT
f.*,
DATE_DIFF(CURRENT_DATE(), f.last_order_date, DAY) as days_since_last_order,
DATE_DIFF(f.last_order_date, f.first_order_date, DAY) as customer_lifetime_days,
COALESCE(t.next_month_value, 0) as target
FROM features f
LEFT JOIN target t ON f.customer_id = t.customer_id
"""
df = self.bq_client.query(query).to_dataframe()
return df
X = df[feature_columns].fillna(0)
y = df['target']
# Split dados
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=0.2, random_state=42
)
# Normalizar features
scaler = StandardScaler()
X_train_scaled = scaler.fit_transform(X_train)
X_test_scaled = [Link](X_test)
# Treinar modelo
model = RandomForestRegressor(
n_estimators=100,
max_depth=10,
random_state=42,
n_jobs=-1
)
[Link](X_train_scaled, y_train)
# Avaliar modelo
train_score = [Link](X_train_scaled, y_train)
test_score = [Link](X_test_scaled, y_test)
# Salvar artefatos
timestamp = [Link]().strftime('%Y%m%d_%H%M%S')
model_path = f"models/{model_name}_{timestamp}"
# Salvar modelo
model_blob = [Link](f"{model_path}/[Link]")
model_blob.upload_from_string([Link](model))
# Salvar scaler
scaler_blob = [Link](f"{model_path}/[Link]")
scaler_blob.upload_from_string([Link](scaler))
# Salvar metadados
metadata = {
'model_name': model_name,
'timestamp': timestamp,
'features': feature_columns,
'train_score': train_score,
'test_score': test_score,
'train_samples': len(X_train),
'test_samples': len(X_test)
}
metadata_blob = [Link](f"{model_path}/[Link]")
metadata_blob.upload_from_string([Link](metadata))
model = [Link](
display_name=model_display_name,
artifact_uri=f"gs://{bucket_name}/{model_path}",
serving_container_image_uri="[Link]/vertex-ai/prediction/sklearn-cpu.1-0:latest"
)
return model
if endpoints:
endpoint = endpoints[0]
else:
endpoint = [Link](
display_name=endpoint_name,
description="Endpoint para previsão de valor do cliente"
)
# Deploy do modelo
deployed_model = [Link](
model=model,
deployed_model_display_name=model.display_name,
machine_type="n1-standard-2",
min_replica_count=1,
max_replica_count=3,
accelerator_type=None,
accelerator_count=0
)
return endpoint
job = model.batch_predict(
job_display_name=f"batch-prediction-{[Link]().strftime('%Y%m%d-%H%M%S')}",
bigquery_source=input_dataset,
bigquery_destination_prefix=output_dataset,
machine_type="n1-standard-4",
starting_replica_count=1,
max_replica_count=5
)
[Link]()
return job
# Uso do pipeline
if __name__ == "__main__":
pipeline = MLPipeline(project_id="seu-projeto-gcp")
# Preparar dados
df = pipeline.prepare_training_data("ecommerce", "orders")
# Treinar modelo
model = pipeline.train_model(df, "customer_ltv_predictor")
# Deploy
endpoint = pipeline.deploy_model(model, "customer-ltv-endpoint")
python
# architecture/data_platform.py
"""
Arquitetura da Plataforma de Dados
Componentes:
1. Ingestão: Cloud Functions + Pub/Sub
2. Processamento: Dataflow + BigQuery
3. Armazenamento: Cloud SQL + Firestore + BigQuery
4. Análise: Vertex AI + BigQuery ML
5. Visualização: Looker Studio + API customizada
6. Monitoramento: Cloud Monitoring + Logging
"""
class DataPlatform:
"""Plataforma integrada de dados em GCP"""
def create_infrastructure(self):
"""Cria toda a infraestrutura necessária"""
try:
dataset = self.bq_client.create_dataset(dataset, timeout=30)
print(f"Dataset criado: {dataset_id}")
except Exception:
print(f"Dataset já existe: {dataset_id}")
# Dataflow Pipeline
class EventProcessor([Link]):
"""Processa eventos em streaming"""
# Enriquecer evento
event['processed_at'] = [Link]().isoformat()
event['processing_version'] = '1.0'
except Exception as e:
# Log de erro
error_event = {
'error': str(e),
'raw_data': [Link]('utf-8', errors='ignore'),
'timestamp': [Link]().isoformat()
}
yield [Link]('errors', error_event)
class AggregateEvents([Link]):
"""Agrega eventos por janela de tempo"""
aggregation = {
'window_start': [Link].to_utc_datetime().isoformat(),
'window_end': [Link].to_utc_datetime().isoformat(),
'event_count': len(events_list),
'unique_users': len(set(e['user_id'] for e in events_list)),
'event_types': {}
}
yield aggregation
pipeline_options = PipelineOptions(
project=project_id,
runner='DataflowRunner',
temp_location=f'gs://{project_id}-dataflow-temp/temp',
region='us-central1',
streaming=True,
save_main_session=True
)
# Schema BigQuery
event_schema = {
'fields': [
{'name': 'event_id', 'type': 'STRING', 'mode': 'REQUIRED'},
{'name': 'user_id', 'type': 'STRING', 'mode': 'REQUIRED'},
{'name': 'event_type', 'type': 'STRING', 'mode': 'REQUIRED'},
{'name': 'timestamp', 'type': 'TIMESTAMP', 'mode': 'REQUIRED'},
{'name': 'properties', 'type': 'JSON', 'mode': 'NULLABLE'},
{'name': 'processed_at', 'type': 'TIMESTAMP', 'mode': 'REQUIRED'}
]
}
aggregation_schema = {
'fields': [
{'name': 'window_start', 'type': 'TIMESTAMP', 'mode': 'REQUIRED'},
{'name': 'window_end', 'type': 'TIMESTAMP', 'mode': 'REQUIRED'},
{'name': 'event_count', 'type': 'INTEGER', 'mode': 'REQUIRED'},
{'name': 'unique_users', 'type': 'INTEGER', 'mode': 'REQUIRED'},
{'name': 'event_types', 'type': 'JSON', 'mode': 'REQUIRED'}
]
}
Conceitos Avançados
Query optimization
Cost management
python
# monitoring/platform_monitor.py
from [Link] import monitoring_v3
from [Link] import logging
from [Link] import bigquery
import pandas as pd
from datetime import datetime, timedelta
from typing import List, Dict, Tuple
import smtplib
from [Link] import MIMEText
class PlatformMonitor:
"""Sistema de monitoramento e otimização da plataforma"""
query = f"""
SELECT
user_email,
DATE(creation_time) as query_date,
query,
total_bytes_processed,
total_slot_ms,
ROUND(total_bytes_processed / POW(10, 12) * 5, 2) as estimated_cost_usd
FROM `{self.project_id}.region-us.INFORMATION_SCHEMA.JOBS_BY_PROJECT`
WHERE creation_time >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL {days} DAY)
AND state = 'DONE'
AND statement_type = 'SELECT'
ORDER BY total_bytes_processed DESC
"""
df = self.bq_client.query(query).to_dataframe()
return df
def optimize_queries(self, expensive_queries: [Link]) -> List[Dict]:
"""Sugere otimizações para queries caras"""
optimizations = []
# Verificar SELECT *
if 'select *' in query_text:
[Link]({
'issue': 'SELECT * detectado',
'suggestion': 'Especifique apenas as colunas necessárias',
'potential_savings': '60-90%'
})
if suggestions:
[Link]({
'query_date': query_info['query_date'],
'user': query_info['user_email'],
'current_cost': query_info['estimated_cost_usd'],
'suggestions': suggestions
})
return optimizations
def create_monitoring_dashboard(self):
"""Cria métricas customizadas e alertas"""
# Métrica customizada para latência de pipeline
descriptor = monitoring_v3.MetricDescriptor(
type=f"[Link]/{self.project_id}/pipeline_latency",
metric_kind=monitoring_v3.[Link],
value_type=monitoring_v3.[Link],
description="Latência do pipeline de dados em segundos",
display_name="Pipeline Latency"
)
try:
self.monitoring_client.create_metric_descriptor(
name=self.project_path,
metric_descriptor=descriptor
)
except Exception:
pass # Métrica já existe
return alert_policy
report = {
'generated_at': [Link]().isoformat(),
'period': 'last_30_days',
'metrics': {}
}
# BigQuery Performance
bq_query = """
SELECT
COUNT(*) as total_queries,
SUM(total_bytes_processed) / POW(10, 12) as tb_processed,
AVG(total_slot_ms) / 1000 as avg_slot_seconds,
SUM(total_bytes_processed) / POW(10, 12) * 5 as total_cost_usd
FROM `region-us.INFORMATION_SCHEMA.JOBS_BY_PROJECT`
WHERE creation_time >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 30 DAY)
AND state = 'DONE'
"""
bq_metrics = self.bq_client.query(bq_query).to_dataframe().iloc[0].to_dict()
report['metrics']['bigquery'] = bq_metrics
# Storage Usage
storage_query = """
SELECT
schema_name as dataset,
SUM(size_bytes) / POW(10, 9) as size_gb
FROM `region-us.INFORMATION_SCHEMA.TABLE_STORAGE`
GROUP BY dataset
ORDER BY size_gb DESC
"""
storage_df = self.bq_client.query(storage_query).to_dataframe()
report['metrics']['storage'] = {
'total_gb': storage_df['size_gb'].sum(),
'by_dataset': storage_df.to_dict('records')
}
# Pipeline Health
recent_errors = self._get_recent_errors()
report['metrics']['pipeline_health'] = {
'error_count': len(recent_errors),
'error_rate': len(recent_errors) / max(bq_metrics['total_queries'], 1),
'top_errors': recent_errors[:5]
}
return report
errors = []
for entry in self.logging_client.list_entries(filter_=filter_str):
[Link]({
'timestamp': [Link](),
'severity': [Link],
'message': [Link]('message', str([Link])),
'resource': [Link]
})
def validate_schema_changes(self,
dataset_id: str,
table_id: str,
new_schema: List[[Link]]) -> Tuple[bool, List[str]]:
"""Valida mudanças de schema"""
table_ref = f"{self.project_id}.{dataset_id}.{table_id}"
try:
table = self.bq_client.get_table(table_ref)
current_schema = [Link]
issues = []
except Exception as e:
# Tabela não existe ainda
return True, []
compatible_changes = {
'INTEGER': ['NUMERIC', 'FLOAT', 'STRING'],
'NUMERIC': ['STRING'],
'FLOAT': ['STRING'],
'STRING': [], # String não pode mudar para outros tipos
'TIMESTAMP': ['STRING'],
'DATE': ['STRING', 'TIMESTAMP'],
'TIME': ['STRING'],
'DATETIME': ['STRING', 'TIMESTAMP']
}
Requisitos
Ingestão em tempo real de múltiplas fontes
Processamento com baixa latência
Monitoramento completo
Implementação
[O código do projeto capstone seria extenso demais para incluir aqui, mas incluiria:]
Recursos Adicionais
Certificações Recomendadas
1. Google Cloud Professional Data Engineer
2. Google Cloud Professional Machine Learning Engineer
Próximos Passos
1. Especializações: Streaming (Apache Beam), ML Ops, Data Mesh
2. Linguagens: Go para performance, Rust para sistemas
DataOps Brasil
MLOps Community
Conclusão
Este curso fornece uma base sólida e prática para engenharia de dados moderna. A progressão foi desenhada
para construir competências incrementalmente, sempre com foco em projetos reais e aplicáveis.
Lembre-se: a excelência vem da prática constante e da curiosidade em explorar novas soluções. Continue
construindo, quebrando e reconstruindo.
"O código é poesia, os dados são a tinta, e a nuvem é nossa tela infinita."