Bir kaydın platform içindeki yolculuğu, her adımın kodu ve verinin o adımdaki hali
… · sayaç …Enerji sayaç okumaları Oracle'da, telekom arama kayıtları PostgreSQL'de tutuluyor. Demo için iki simülatör bu tablolara 10 saniyede bir kayıt ekliyor (%4'ü kasıtlı olarak bozuk). Enerji kesintileri ve SMS olayları ise Kafka'ya JSON mesaj olarak akıyor.
CREATE TABLE EDP_DEMO.METER_READINGS (
reading_id VARCHAR2(36) PRIMARY KEY,
meter_id VARCHAR2(32) NOT NULL,
reading_ts TIMESTAMP NOT NULL,
kwh NUMBER(12,3) NOT NULL,
voltage NUMBER(8,2) NOT NULL,
region VARCHAR2(32),
city VARCHAR2(32),
tariff VARCHAR2(16),
power_factor NUMBER(4,3),
status VARCHAR2(16)
-- + district, substation, feeder_id, phase,
-- reactive_kvarh, current_a, frequency_hz, temperature_c
);INSERT INTO EDP_DEMO.METER_READINGS (
reading_id, meter_id, reading_ts, kwh, voltage,
region, city, district, substation, feeder_id, tariff, phase,
reactive_kvarh, current_a, power_factor, frequency_hz,
temperature_c, status
) VALUES (
:reading_id, :meter_id, :reading_ts, :kwh, :voltage,
:region, :city, :district, :substation, :feeder_id, :tariff, :phase,
:reactive_kvarh, :current_a, :power_factor, :frequency_hz,
:temperature_c, :status
)CREATE TABLE IF NOT EXISTS edp_demo.cdr_events (
event_id uuid PRIMARY KEY,
msisdn varchar(20) NOT NULL,
event_ts timestamptz NOT NULL,
duration_sec integer NOT NULL,
cell_id varchar(32) NOT NULL
);
INSERT INTO edp_demo.cdr_events
(event_id, msisdn, event_ts, duration_sec, cell_id)
VALUES (%s, %s, %s, %s, %s);{
"message_id": "0750408b-1070-4675-8e2b-2e9ab1c6fffe",
"msisdn": "905551002610",
"event_ts": "2026-09-24T06:44:16Z",
"cell_id": "CELL-14A",
"direction": "mo"
}
// direction: mo = giden, mt = gelen
// her 9. mesajda "xx" (bozuk) -> silver'da elenirHer kaynak için bir SeaTunnel job'u var. Job iki bloktan oluşur: source nereden okunacağını (SELECT), sink nereye yazılacağını söyler. INSERT yazılmaz; engine satırları okuyup hedefe kendisi yazar. Burada hedef MinIO'daki landing alanı (Parquet / JSON dosyaları).
env { job.mode = "BATCH", parallelism = 1,
job.name = "edp_ingest_energy_meter_readings" }
source {
Jdbc {
url = "jdbc:oracle:thin:@//$ENERGY_DB_DSN"
driver = "oracle.jdbc.OracleDriver"
user = "$ENERGY_DB_USER"
password = "$ORACLE_PASSWORD"
query = "SELECT READING_ID AS reading_id,
METER_ID AS meter_id, READING_TS AS reading_ts,
KWH AS kwh, VOLTAGE AS voltage, ...
FROM EDP_DEMO.METER_READINGS"
}
}sink {
S3File {
bucket = "s3a://iceberg-warehouse"
path = "/edp/landing/energy_meter_readings"
fs.s3a.endpoint = "http://minio.lakehouse.svc:9000"
file_format_type = "parquet" # Kafka icin: json
data_save_mode = "DROP_DATA" # once eskiyi sil, sonra yaz
schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
}
}
# Kafka job'larinda source:
# Kafka { topic = "telecom.sms", start_mode = "earliest",
# format = json, format_error_handle_way = "skip" }Veri kaynaktaki haliyle, hiç değiştirilmeden duruyor. Oracle'ın NUMBER alanları burada hâlâ decimal, Kafka'dan gelen tarih hâlâ metin. Bozuk kayıtlar da burada, çünkü landing'in görevi "geleni olduğu gibi kabul etmek". Sorgulanabilmesi için aynı veri Iceberg tablosu olarak da tutuluyor (iceberg.landing_zone.*).
Spark landing dosyalarını okuyor, veri tiplerini standartlaştırıyor ve her kayda nereden geldiğini (source_system) ve ne zaman yüklendiğini (ingested_at) ekliyor. Yazma işi tek satır: writeTo(...).createOrReplace(). Bu, SQL'deki CREATE OR REPLACE TABLE ... AS SELECT ile aynı iş. Tablo formatı Apache Iceberg, katalog Apache Polaris.
spark = (SparkSession.builder.appName("edp-bronze")
.config("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg.type", "rest")
.config("spark.sql.catalog.iceberg.uri", "http://polaris.lakehouse.svc:8181/api/catalog")
.config("spark.sql.catalog.iceberg.warehouse", "edp")
.getOrCreate())
def landed(name, fmt="parquet"):
df = spark.read.parquet(f"s3a://iceberg-warehouse/edp/landing/{name}")
# ham kopya -> iceberg.landing_zone (sorgulanabilir landing)
df.withColumn("_loaded_at", current_timestamp()) \
.writeTo(f"iceberg.landing_zone.{name}").createOrReplace()
return df
def save(df, table):
df.writeTo(f"iceberg.bronze.{table}").createOrReplace() # <- "INSERT"
meter_raw = landed("energy_meter_readings")
save(meter_raw.select(
col("reading_id").cast("string"),
col("reading_ts").cast("timestamp"),
col("kwh").cast("double"), # decimal(12,3) -> double
col("voltage").cast("double"),
...,
lit("oracle").alias("source_system"),
current_timestamp().alias("ingested_at")),
"energy_meter_readings")Aynı kayıt bronze'da: değerler aynı ama tipler standart (kwh artık double) ve izleme kolonları eklenmiş. Bozuk kayıtlar hâlâ burada. Bronze, "kaynağın izlenebilir, yeniden işlenebilir kopyası".
dbt modeli sadece bir SELECT: hangi kalite kuralları uygulanacak ve tekrar eden kayıtlar nasıl temizlenecek. materialized: table ayarı sayesinde dbt bu SELECT'i CREATE OR REPLACE TABLE komutuna sarıp Trino'ya gönderiyor. Hesaplama Trino'da, veri yine Iceberg/MinIO'da.
select reading_id, meter_id, reading_ts, region, city,
tariff, kwh, voltage, power_factor, status,
source_system, ingested_at
from (
select *,
row_number() over (partition by reading_id
order by ingested_at desc) as rn
from {{ source('bronze', 'energy_meter_readings') }}
where kwh > 0 -- kalite kurallari
and voltage between 200 and 260
and meter_id is not null
and (power_factor is null or power_factor between 0.7 and 1.0)
and (frequency_hz is null or frequency_hz between 49 and 51)
and (status is null or status <> 'bad')
) deduped
where rn = 1 -- tekillestirme/* {"app": "dbt", "node_id": "model.edp.energy_meter_readings"} */
create or replace table "iceberg"."silver"."energy_meter_readings" as (
select reading_id, meter_id, reading_ts, ...
from (
select *, row_number() over (...) as rn
from "iceberg"."bronze"."energy_meter_readings"
where kwh > 0 and voltage between 200 and 260 ...
) deduped
where rn = 1
);Takip ettiğimiz kayıt kuralları geçti ve silver'da. Kuralları geçemeyen kayıtlar bronze'da kalıyor ama silver'a girmiyor; aşağıda son elenen kaydı ve neden elendiğini görüyorsun.
Gold, raporlamaya hazır veri ürünleri. Takip ettiğimiz okuma, sayacın o saatlik özetine katkı veriyor: aynı saatteki diğer okumalarla birlikte toplanıp tek satıra dönüşüyor.
select meter_id, region, city, feeder_id, tariff,
date_trunc('hour', reading_ts) as hour_ts,
count(*) as reading_count,
round(sum(kwh), 3) as total_kwh,
round(sum(reactive_kvarh), 3) as total_reactive_kvarh,
round(avg(voltage), 2) as avg_voltage,
round(avg(current_a), 2) as avg_current_a,
round(avg(power_factor), 3) as avg_power_factor
from {{ ref('energy_meter_readings') }} -- silver modeli
group by 1, 2, 3, 4, 5, 6Gold tablolarına standart SQL ile erişiliyor. Power BI, Excel ya da herhangi bir uygulama Trino'ya JDBC/ODBC ile bağlanıp aynı sorguyu çalıştırabilir. Iceberg her yazımı ayrı bir sürüm (snapshot) olarak sakladığı için geçmişe dönük sorgu da mümkün.
-- Son 6 saatin bolge bazinda enerji tuketimi
SELECT hour_ts AS saat,
region AS bolge,
SUM(reading_count) AS okuma,
ROUND(SUM(total_kwh), 1) AS toplam_kwh
FROM iceberg.gold.energy_meter_hourly
WHERE hour_ts >= date_trunc('hour', current_timestamp)
- INTERVAL '6' HOUR
GROUP BY 1, 2
ORDER BY saat DESC, bolge;-- 30 dakika once kac kayit vardi, simdi kac?
SELECT
(SELECT COUNT(*) FROM iceberg.bronze.energy_meter_readings
FOR TIMESTAMP AS OF current_timestamp
- INTERVAL '30' MINUTE) AS otuz_dk_once,
(SELECT COUNT(*)
FROM iceberg.bronze.energy_meter_readings) AS simdi;Tüm adımları tek bir DAG sırayla çalıştırıyor. Her adım Kubernetes üzerinde ayrı bir pod olarak açılıp işi bitince kapanıyor; bir adım hata verirse sonrakiler çalışmıyor. Airflow orkestrasyonu yapıyor, hesaplama SeaTunnel, Spark ve Trino'da.
with DAG(
dag_id="edp_pipeline",
schedule=timedelta(minutes=10),
max_active_runs=1, # calisma bitmeden yenisi baslamaz
default_args={"retries": 1},
catchup=False,
):
ingest = seatunnel_task("seatunnel_ingest")
bronze = spark_task("spark_bronze", "run_bronze.py", ...)
silver = dbt_task("dbt_silver", "tag:silver")
gold = dbt_task("dbt_gold", "tag:gold")
ingest >> bronze >> silver >> gold