Untitled
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