← Mimari sayfasına dön

Veri Akışı: Kaynaktan Gold'a

Bir kaydın platform içindeki yolculuğu, her adımın kodu ve verinin o adımdaki hali

smarttech
Canlı örnek veri: yükleniyor…
Kaynak sistemlere sürekli yeni kayıt yazılıyor. Airflow edp_pipeline her 10 dakikada bir çalışıyor: SeaTunnel veriyi kaynaklardan alıp MinIO'ya (landing) bırakıyor, Spark bronze katmanına yazıyor, dbt Trino üzerinde silver ve gold tablolarını üretiyor. Aşağıda aynı sayaç okumasını katman katman takip ediyoruz; veri tabloları her 10 dakikada gerçek sistemden yenileniyor.
Takip edilen kayıt: … · sayaç …
0

Kaynak sistemler

OraclePostgreSQLKafka SQL Editör'de aç ↗

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.

Oracle · EDP_DEMO.METER_READINGSDDL
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
);
energy-oracle-sim · sim.py10 sn'de bir
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
)
PostgreSQL · edp_demo.cdr_eventsDDL + INSERT
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);
Kafka · topic telecom.smsJSON mesaj
{
  "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 elenir
1

Veri alımı: Apache SeaTunnel

ingestbatch SeaTunnel'i aç ↗

Her 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ı).

seatunnel/ingest.sh · Oracle jobsource
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"
  }
}
seatunnel/ingest.sh · ortak hedefsink = "INSERT"
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" }
Şifreler koda yazılmıyor; Kubernetes secret'ından ortam değişkeni olarak geliyor ($ORACLE_PASSWORD gibi). Aynı değerler Vault'ta da secret/edp/source/* altında duruyor.
2

Landing

MinIO · ham veri MinIO'da gör ↗

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.*).

Takip edilen kayıt · landing_zone.energy_meter_readings CANLI

Kafka'dan bir SMS olayı · landing_zone.telecom_sms CANLI

3

Dönüşüm: Apache Spark

landing → bronze Airflow'da task ↗

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/run_bronze.pyPySpark
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")
4

Bronze

Iceberg · iceberg.bronze SQL: 04_Bronze ↗

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ı".

Takip edilen kayıt · bronze.energy_meter_readings CANLI

5

Dönüşüm: dbt (Trino üzerinde)

bronze → silver dbt Docs'ta modeli aç ↗

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.

dbt/models/silver/energy_meter_readings.sqldbt model
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
dbt'nin Trino'ya gönderdiği SQLTrino sorgu geçmişi
/* {"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
);
dbt'nin gönderdiği gerçek komutlar Trino arayüzünde kullanıcısı dbt olan sorgular olarak görünüyor. Her çalışmada 4 silver + 4 gold tablosu.
6

Silver

Iceberg · iceberg.silver SQL: 05_Silver ↗

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.

Takip edilen kayıt · silver.energy_meter_readings CANLI

Elenen son kayıt (bronze'da var, silver'da yok) CANLI

7

Gold

Iceberg · iceberg.gold dbt Docs: lineage ↗

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.

dbt/models/gold/energy_meter_hourly.sqldbt model
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, 6

Takip edilen sayacın o saatlik satırı · gold.energy_meter_hourly CANLI

8

Tüketim: Trino SQL

BI · analitik · uygulamalar SQL Editör ↗

Gold 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.

SQL Editör · 06_Gold.sqlTrino
-- 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;
SQL Editör · 07_Pipeline_Izleme.sqlIceberg time travel
-- 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;
9

Orkestrasyon: Apache Airflow

her 10 dakika edp_pipeline ↗

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.

airflow/dags/edp_pipeline.pyDAG
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

Katman katman kayıt sayısı (sayaç okumaları) CANLI

Son pipeline çalışmaları (bronze tablosu snapshot'ları) CANLI