Untitled

 avatar
alp
plain_text
9 months ago
18 kB
12
Indexable
"""
CDC Entegrasyonu ve Otomatik Tablo Bakımı
- CDC ile gelen verilerin düzgün yazılması
- Günlük/haftalık otomatik optimizasyon
- Performance monitoring
"""

from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from datetime import datetime, timedelta
import time

# ============================================================================
# KONFIGURASYON
# ============================================================================

TABLE_NAME = "iceberg.sompodatabase2.T001PSRCVP"
DATE_COL = "CONFIRM_DATE"
PK_COLS = ["FIRM_CODE", "COMPANY_CODE", "PRODUCT_NO", "POLICY_NO", "RENEWAL_NO", "ENDORS_NO"]

# Bakım parametreleri
OPTIMIZE_THRESHOLD_FILES = 50  # Bu kadar dosya birikirse optimize et
OPTIMIZE_THRESHOLD_SIZE_MB = 50  # Ortalama dosya boyutu bu değerin altındaysa optimize et
TARGET_FILE_SIZE_MB = 512
SNAPSHOT_RETENTION_DAYS = 7
ORPHAN_FILE_RETENTION_DAYS = 3

# ============================================================================
# 1. CDC WRITE KONFIGURASYONU
# ============================================================================

def configure_cdc_write_properties(spark, table_name):
    """
    CDC için optimal write ayarlarını yapılandır
    Bu ayarları CDC pipeline'ınızda kullanın
    """
    print("="*80)
    print("CDC WRITE KONFIGÜRASYONU")
    print("="*80)
    
    # Tablo özelliklerini ayarla
    spark.sql(f"""
        ALTER TABLE {table_name}
        SET TBLPROPERTIES (
            -- Write performansı
            'write.format.default' = 'parquet',
            'write.parquet.compression-codec' = 'zstd',
            'write.target-file-size-bytes' = '{TARGET_FILE_SIZE_MB * 1024 * 1024}',
            
            -- Upsert/Merge için optimize
            'write.merge.mode' = 'copy-on-write',
            'write.update.mode' = 'copy-on-write',
            'write.delete.mode' = 'copy-on-write',
            
            -- Distribution
            'write.distribution-mode' = 'hash',
            
            -- Metadata
            'write.metadata.compression-codec' = 'gzip',
            'write.metadata.metrics.default' = 'full',
            
            -- Commit
            'commit.retry.num-retries' = '10',
            'commit.retry.min-wait-ms' = '100',
            'commit.manifest.min-count-to-merge' = '5',
            
            -- History
            'history.expire.max-snapshot-age-ms' = '{SNAPSHOT_RETENTION_DAYS * 24 * 3600 * 1000}'
        )
    """)
    
    print(f"✓ Tablo özellikleri CDC için optimize edildi")
    print(f"""
CDC Pipeline'ınızda şu Spark config'leri kullanın:

spark.conf.set("spark.sql.iceberg.planning-mode", "distributed")
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.files.maxPartitionBytes", "536870912")  # 512MB

# CDC Write örneği:
df.writeTo("{table_name}") \\
  .option("write-format", "parquet") \\
  .option("compression", "zstd") \\
  .partitionedBy("months(CONFIRM_DATE)") \\
  .append()  # veya .createOrReplace() / .merge()
    """)

# ============================================================================
# 2. PARTITION SAĞLIK KONTROLÜ
# ============================================================================

def check_partition_health(spark, table_name, date_col):
    """
    Hangi partition'ların optimize edilmesi gerektiğini tespit et
    """
    print("\n" + "="*80)
    print("PARTITION SAĞLIK KONTROLÜ")
    print("="*80)
    
    # Son 30 günün partition'larını kontrol et
    unhealthy_partitions = spark.sql(f"""
        SELECT 
            DATE_TRUNC('MONTH', {date_col}) as partition_month,
            COUNT(*) as file_count,
            ROUND(SUM(file_size_in_bytes) / 1024 / 1024 / 1024, 2) as size_gb,
            ROUND(AVG(file_size_in_bytes) / 1024 / 1024, 2) as avg_file_mb,
            ROUND(MIN(file_size_in_bytes) / 1024 / 1024, 2) as min_file_mb,
            ROUND(MAX(file_size_in_bytes) / 1024 / 1024, 2) as max_file_mb
        FROM {table_name}.files
        WHERE {date_col} >= ADD_MONTHS(CURRENT_DATE(), -1)
        GROUP BY DATE_TRUNC('MONTH', {date_col})
        HAVING 
            COUNT(*) > {OPTIMIZE_THRESHOLD_FILES} 
            OR AVG(file_size_in_bytes) < {OPTIMIZE_THRESHOLD_SIZE_MB * 1024 * 1024}
        ORDER BY partition_month DESC
    """)
    
    unhealthy_list = unhealthy_partitions.collect()
    
    if len(unhealthy_list) > 0:
        print(f"⚠️  {len(unhealthy_list)} partition optimize edilmeli:")
        unhealthy_partitions.show(truncate=False)
        return [row['partition_month'] for row in unhealthy_list]
    else:
        print("✓ Tüm partition'lar sağlıklı durumda")
        return []

# ============================================================================
# 3. GÜNLÜK OTOMATİK OPTİMİZASYON
# ============================================================================

def daily_optimization(spark, table_name, date_col, pk_cols):
    """
    Günlük çalıştırılacak optimizasyon görevi
    - CDC ile değişen partition'ları optimize et
    - Küçük dosyaları birleştir
    """
    print("\n" + "="*80)
    print(f"GÜNLÜK OPTİMİZASYON - {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
    print("="*80)
    
    start_time = time.time()
    
    # 1. Optimize edilmesi gereken partition'ları tespit et
    partitions_to_optimize = check_partition_health(spark, table_name, date_col)
    
    if not partitions_to_optimize:
        print("✓ Optimize edilecek partition yok")
        return
    
    # 2. Her partition'ı optimize et
    sort_order = ",".join(pk_cols)
    
    for partition_month in partitions_to_optimize:
        month_str = partition_month.strftime('%Y-%m')
        print(f"\n📊 {month_str} optimize ediliyor...")
        
        try:
            result = spark.sql(f"""
                CALL iceberg.system.rewrite_data_files(
                    table => '{table_name}',
                    strategy => 'sort',
                    sort_order => '{sort_order}',
                    where => "{date_col} >= DATE'{partition_month}' AND {date_col} < DATE'{partition_month}' + INTERVAL 1 MONTH",
                    options => map(
                        'target-file-size-bytes', '{TARGET_FILE_SIZE_MB * 1024 * 1024}',
                        'partial-progress.enabled', 'true',
                        'max-concurrent-file-group-rewrites', '8'
                    )
                )
            """).first()
            
            print(f"  ✓ Yeniden yazılan dosya: {result['rewritten_data_files_count']}")
            print(f"  ✓ Yeniden yazılan veri: {result['rewritten_bytes_count'] / 1024 / 1024 / 1024:.2f} GB")
            
        except Exception as e:
            print(f"  ✗ HATA: {str(e)}")
            continue
    
    elapsed = time.time() - start_time
    print(f"\n✓ Günlük optimizasyon tamamlandı ({elapsed/60:.1f} dakika)")

# ============================================================================
# 4. HAFTALIK BAKIM
# ============================================================================

def weekly_maintenance(spark, table_name):
    """
    Haftalık çalıştırılacak bakım görevi
    - Manifest optimize
    - Eski snapshot'ları temizle
    - Orphan dosyaları temizle
    """
    print("\n" + "="*80)
    print(f"HAFTALIK BAKIM - {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
    print("="*80)
    
    # 1. Manifest optimize
    print("\n1/3 Manifest dosyaları optimize ediliyor...")
    try:
        result = spark.sql(f"""
            CALL iceberg.system.rewrite_manifests(
                table => '{table_name}'
            )
        """).first()
        print(f"  ✓ Yeniden yazılan manifest: {result['rewritten_manifests_count']}")
    except Exception as e:
        print(f"  ✗ HATA: {str(e)}")
    
    # 2. Eski snapshot'ları temizle
    print("\n2/3 Eski snapshot'lar temizleniyor...")
    try:
        result = spark.sql(f"""
            CALL iceberg.system.expire_snapshots(
                table => '{table_name}',
                older_than => TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL {SNAPSHOT_RETENTION_DAYS} DAY),
                retain_last => 5
            )
        """).first()
        print(f"  ✓ Silinen snapshot: {result.get('deleted_snapshots_count', 0)}")
    except Exception as e:
        print(f"  ✗ HATA: {str(e)}")
    
    # 3. Orphan dosyaları temizle
    print("\n3/3 Orphan dosyalar temizleniyor...")
    try:
        result = spark.sql(f"""
            CALL iceberg.system.remove_orphan_files(
                table => '{table_name}',
                older_than => TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL {ORPHAN_FILE_RETENTION_DAYS} DAY)
            )
        """).first()
        orphan_count = result.get('orphan_file_location', [])
        print(f"  ✓ Silinen orphan dosya: {len(orphan_count) if orphan_count else 0}")
    except Exception as e:
        print(f"  ✗ HATA: {str(e)}")
    
    print("\n✓ Haftalık bakım tamamlandı")

# ============================================================================
# 5. PERFORMANS İZLEME
# ============================================================================

def monitor_performance(spark, table_name, date_col):
    """
    Tablo performans metriklerini izle
    """
    print("\n" + "="*80)
    print("PERFORMANS METRİKLERİ")
    print("="*80)
    
    # 1. Genel durum
    print("\n1. GENEL DURUM:")
    general_stats = spark.sql(f"""
        SELECT 
            COUNT(*) as total_files,
            ROUND(SUM(file_size_in_bytes) / 1024 / 1024 / 1024, 2) as total_gb,
            ROUND(AVG(file_size_in_bytes) / 1024 / 1024, 2) as avg_file_mb,
            COUNT(DISTINCT partition) as partition_count
        FROM {table_name}.files
    """).first()
    
    print(f"  Toplam dosya: {general_stats['total_files']}")
    print(f"  Toplam boyut: {general_stats['total_gb']} GB")
    print(f"  Ortalama dosya: {general_stats['avg_file_mb']} MB")
    print(f"  Partition sayısı: {general_stats['partition_count']}")
    
    # 2. Son 7 günün partition durumu
    print("\n2. SON 7 GÜNÜN PARTITION DURUMU:")
    recent_partitions = spark.sql(f"""
        SELECT 
            DATE_TRUNC('DAY', {date_col}) as day,
            COUNT(*) as files,
            ROUND(SUM(file_size_in_bytes) / 1024 / 1024, 2) as size_mb,
            ROUND(AVG(file_size_in_bytes) / 1024 / 1024, 2) as avg_mb
        FROM {table_name}.files
        WHERE {date_col} >= DATE_SUB(CURRENT_DATE(), 7)
        GROUP BY DATE_TRUNC('DAY', {date_col})
        ORDER BY day DESC
    """)
    recent_partitions.show(truncate=False)
    
    # 3. Snapshot geçmişi
    print("\n3. SNAPSHOT GEÇMİŞİ (Son 10):")
    snapshots = spark.sql(f"""
        SELECT 
            committed_at,
            snapshot_id,
            operation,
            summary
        FROM {table_name}.snapshots
        ORDER BY committed_at DESC
        LIMIT 10
    """)
    snapshots.show(truncate=False)
    
    # 4. Sağlık skoru
    avg_file_size_mb = general_stats['avg_file_mb']
    file_count = general_stats['total_files']
    
    health_score = 100
    if avg_file_size_mb < 100:
        health_score -= 30
    if avg_file_size_mb < 50:
        health_score -= 20
    if file_count > 1000:
        health_score -= 20
    if file_count > 5000:
        health_score -= 20
    
    print(f"\n4. SAĞLIK SKORU: {health_score}/100")
    if health_score >= 80:
        print("  ✅ Mükemmel durumda")
    elif health_score >= 60:
        print("  ⚠️  İyileştirme gerekebilir")
    else:
        print("  ❌ Acil optimizasyon gerekli")

# ============================================================================
# 6. CDC MERGE ÖRNEĞİ
# ============================================================================

def cdc_merge_example(spark, table_name, cdc_df, pk_cols):
    """
    CDC verilerini merge etme örneği
    
    Parameters:
    - cdc_df: CDC'den gelen DataFrame (op: I=Insert, U=Update, D=Delete)
    """
    print("\n" + "="*80)
    print("CDC MERGE İŞLEMİ")
    print("="*80)
    
    # Temporary view oluştur
    cdc_df.createOrReplaceTempView("cdc_changes")
    
    # Match condition
    match_condition = " AND ".join([f"target.{col} = source.{col}" for col in pk_cols])
    
    # Merge query
    merge_query = f"""
        MERGE INTO {table_name} target
        USING cdc_changes source
        ON {match_condition}
        
        WHEN MATCHED AND source.op = 'D' THEN DELETE
        
        WHEN MATCHED AND source.op = 'U' THEN UPDATE SET *
        
        WHEN NOT MATCHED AND source.op IN ('I', 'U') THEN INSERT *
    """
    
    print("Merge sorgusu çalıştırılıyor...")
    start_time = time.time()
    
    spark.sql(merge_query)
    
    elapsed = time.time() - start_time
    print(f"✓ Merge tamamlandı ({elapsed:.2f} saniye)")
    
    # İşlem sonrası istatistikler
    stats = spark.sql(f"SELECT * FROM {table_name}.snapshots ORDER BY committed_at DESC LIMIT 1").first()
    print(f"  Added data files: {stats['summary'].get('added-data-files', 0)}")
    print(f"  Deleted data files: {stats['summary'].get('deleted-data-files', 0)}")
    print(f"  Added records: {stats['summary'].get('added-records', 0)}")
    print(f"  Deleted records: {stats['summary'].get('deleted-records', 0)}")

# ============================================================================
# 7. AİRFLOW/CRON SCHEDULER ÖRNEĞİ
# ============================================================================

def create_airflow_dag_example():
    """
    Airflow DAG örneği (sadece template)
    """
    dag_code = '''
from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'data-team',
    'depends_on_past': False,
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 3,
    'retry_delay': timedelta(minutes=5),
}

# GÜNLÜK OPTIMIZASYON DAG
with DAG(
    'iceberg_daily_optimization',
    default_args=default_args,
    description='Iceberg tablo günlük optimizasyon',
    schedule_interval='0 2 * * *',  # Her gün 02:00
    start_date=datetime(2025, 1, 1),
    catchup=False,
    tags=['iceberg', 'optimization'],
) as daily_dag:

    optimize_task = SparkSubmitOperator(
        task_id='optimize_partitions',
        application='/path/to/daily_optimization.py',
        conn_id='spark_default',
        conf={
            'spark.sql.adaptive.enabled': 'true',
            'spark.sql.iceberg.planning-mode': 'distributed',
        }
    )

# HAFTALIK BAKIM DAG
with DAG(
    'iceberg_weekly_maintenance',
    default_args=default_args,
    description='Iceberg tablo haftalık bakım',
    schedule_interval='0 3 * * 0',  # Her Pazar 03:00
    start_date=datetime(2025, 1, 1),
    catchup=False,
    tags=['iceberg', 'maintenance'],
) as weekly_dag:

    manifest_task = SparkSubmitOperator(
        task_id='optimize_manifests',
        application='/path/to/weekly_maintenance.py',
        conn_id='spark_default',
    )
    '''
    
    print("\n" + "="*80)
    print("AIRFLOW DAG ÖRNEK KODU")
    print("="*80)
    print(dag_code)

# ============================================================================
# ANA FONKSİYONLAR
# ============================================================================

def setup_cdc_integration():
    """CDC entegrasyonunu kur"""
    spark = SparkSession.builder.getOrCreate()
    configure_cdc_write_properties(spark, TABLE_NAME)
    print("\n✓ CDC entegrasyonu hazır!")

def run_daily_job():
    """Günlük job"""
    spark = SparkSession.builder.getOrCreate()
    daily_optimization(spark, TABLE_NAME, DATE_COL, PK_COLS)
    monitor_performance(spark, TABLE_NAME, DATE_COL)

def run_weekly_job():
    """Haftalık job"""
    spark = SparkSession.builder.getOrCreate()
    weekly_maintenance(spark, TABLE_NAME)
    monitor_performance(spark, TABLE_NAME, DATE_COL)

# ============================================================================
# KULLANIM ÖRNEKLERİ
# ============================================================================

if __name__ == "__main__":
    print("""
    
KULLANIM:

1. CDC Entegrasyonu Kurulumu (bir kere):
   python cdc_maintenance.py --setup

2. Günlük Optimizasyon (cronjob/airflow):
   python cdc_maintenance.py --daily

3. Haftalık Bakım (cronjob/airflow):
   python cdc_maintenance.py --weekly

4. Performans İzleme:
   python cdc_maintenance.py --monitor

5. Airflow DAG Örneği:
   python cdc_maintenance.py --airflow-example

    """)
    
    import sys
    if len(sys.argv) > 1:
        if sys.argv[1] == '--setup':
            setup_cdc_integration()
        elif sys.argv[1] == '--daily':
            run_daily_job()
        elif sys.argv[1] == '--weekly':
            run_weekly_job()
        elif sys.argv[1] == '--monitor':
            spark = SparkSession.builder.getOrCreate()
            monitor_performance(spark, TABLE_NAME, DATE_COL)
        elif sys.argv[1] == '--airflow-example':
            create_airflow_dag_example()
Editor is loading...
Leave a Comment