26Y06d

Young-Kyoo Kim·2026년 7월 6일
#!/bin/bash
# =========================================================================
# run_pipeline_analysis.sh — Granular Date-by-Date FinOps Analysis Pipeline
# =========================================================================
set -e

echo "⚙️  [컨테이너 가동] 쿠버네티스 인젝션 환경변수 검증 및 일자별 분할 통제 개시..."
echo "• TARGET_CLUSTER : ${CLUSTER}"
echo "• RELATIVE_DAYS  : ${DAYS:-NOT_SET}"
echo "• ABSOLUTE_RANGE : ${START_DATE:-NOT_SET} ~ ${END_DATE:-NOT_SET}"

# 🛡️ 필수 환경변수 가드레일
if [ -z "$CLUSTER" ]; then
  echo "❌ [에러] CLUSTER 환경변수가 누락되었습니다. (COMPUTE 또는 STORAGE 필수)"
  exit 1
fi

if [ -z "$DAYS" ] && { [ -z "$START_DATE" ] || [ -z "$END_DATE" ]; }; then
  echo "❌ [에러] 기간 설정이 누락되었습니다. DAYS 또는 START_DATE & END_DATE 세트가 필요합니다."
  exit 1
fi

# ─── 🐍 [일자별 스마트 오케스트레이션 엔진] ───
# 파이썬 내부에서 일자별 루프를 돌며 AIStor의 3대 가공원부 경로를 디버깅 출력 및 검증합니다.
python3 -c "
import os, sys, boto3, subprocess
from datetime import datetime, timedelta, timezone

KST = timezone(timedelta(hours=9))
cluster = os.getenv('CLUSTER', '').upper()
days_str = os.getenv('DAYS', '')
start_date = os.getenv('START_DATE', '')
end_date = os.getenv('END_DATE', '')

s3_client = boto3.client(
    's3', 
    endpoint_url=os.getenv('MINIO_ENDPOINT'),
    aws_access_key_id=os.getenv('MINIO_ACCESS_KEY'),
    aws_secret_access_key=os.getenv('MINIO_SECRET_KEY')
)
bucket_name = os.getenv('MINIO_RAW_BUCKET', 'devops-test')
now_kst = datetime.now(KST)

# 1. 분석 대상 날짜 배열 파싱 (상대/절대 기간 모두 대응)
target_dates = []
if start_date and end_date:
    try:
        s_dt = datetime.strptime(start_date, '%Y-%m-%d')
        e_dt = datetime.strptime(end_date, '%Y-%m-%d')
        curr = s_dt
        while curr <= e_dt:
            target_dates.append(curr)
            curr += timedelta(days=1)
    except Exception as e:
        print(f'❌ 날짜 포맷 파싱 에러: {str(e)}')
        sys.exit(1)
elif days_str:
    for i in range(1, int(days_str) + 1):
        target_dates.append(now_kst - timedelta(days=i))

# 로컬 캐시 디렉토리 보장
local_dir = f'./data/merged/{cluster}'
os.makedirs(local_dir, exist_ok=True)

print(f'\n🔍 총 {len(target_dates)}개 일자에 대한 파티션 원부 정밀 검증 세션 가동...\n')

# 2. 📅 일자별 독립 검증 및 분기 루프 집행
for dt in target_dates:
    date_str = dt.strftime('%Y%m%d')
    date_hyphen = dt.strftime('%Y-%m-%d')
    
    # 검증 대상 3대 마스터 파일 바인딩
    f_enriched = f'daily_enriched_{cluster}_{date_str}.parquet'
    f_ns_usage = f'daily_ns_usage_{cluster}_{date_str}.parquet'
    f_pareto   = f'pareto_ns_{cluster}_{date_str}.parquet'
    
    artifact_files = [f_enriched, f_ns_usage, f_pareto]
    all_exist = True
    
    print(f'===================== 📅 TARGET DATE: {date_hyphen} =====================')
    
    # 🔍 [디버그 레이어] 각 파일별 AIStor 물리 경로 로그 출력 및 존재 여부 조사
    for f in artifact_files:
        object_key = f'merged/{cluster}/{f}'
        print(f'[DEBUG] Checking AIStor Path ➡️ s3://{bucket_name}/{object_key}')
        try:
            s3_client.head_object(Bucket=bucket_name, Key=object_key)
        except Exception:
            print(f'[DEBUG]   ❌ 미검출 (유실 또는 미정산 자산): {f}')
            all_exist = False

    # 🔀 분할 정복 조건부 라우팅
    if all_exist:
        # ⚡ [CACHE HIT] 3개 파일이 모두 스토리지에 있으므로 로컬 디스크로 바로 땡겨옴 (다운로드)
        print(f'⚡ [DEBUG] [CACHE HIT] {date_hyphen} 정산본 발견. AIStor 레이어에서 다이렉트 다운로드 진행...')
        for f in artifact_files:
            object_key = f'merged/{cluster}/{f}'
            local_path = os.path.join(local_dir, f)
            print(f'[DEBUG]   📥 Fetching: s3://{bucket_name}/{object_key} ➡️ {local_path}')
            s3_client.download_file(bucket_name, object_key, local_path)
    else:
        # 💨 [CACHE MISS] 가공 파일이 유실되었으므로 해당 일자만 지정해 2단계(Polars) 강제 구동
        print(f'💨 [DEBUG] [CACHE MISS] {date_hyphen} 정산본 유실 감지! Raw 데이터 기반 단일 정산 마감 트리거...')
        
        # 1단계 수집 없이 기존에 AIStor에 적재되어 있던 raw/.parquet를 2단계 스마트 캐시가 알아서 땡겨와 연산합니다.
        cmd = [
            'python3', 'step2_pipeline_polars.py',
            '--cluster', cluster,
            '--start-date', date_hyphen,
            '--end-date', date_hyphen
        ]
        print(f'[DEBUG]   🏃 Executing: {\" \".join(cmd)}')
        
        # 실시간 서브프로세스 파이프 연결 및 가동
        res = subprocess.run(cmd, capture_output=False, text=True)
        if res.returncode != 0:
            print(f'❌ [장애 발생] {date_hyphen} 일자 리커버리 정산 중 2단계 엔진이 크래시되었습니다.')
            sys.exit(1)

print('\n✅ [디버그 완료] 모든 분석 대상 일자의 가공 원부 정렬 및 로컬 동기화 완수.')
"

# ─── 🚀 후속 분석 레이어 연속 집행 ───
# 위 파이썬 루프가 끝나면 로컬의 data/merged/{CLUSTER}/ 하방에 대상 일자의 3대 파일이 무결하게 정착됩니다.
echo ""
echo "========================================================================="
echo "📊 후속 거버넌스 파이프라인 가동 (3단계 차트 ➡️ 4단계 분석 ➡️ 6단계 마스터 엑셀)"
echo "========================================================================="

python3 step3_analytics.py --cluster "$CLUSTER"
python3 step4_governance_analyzer.py --cluster "$CLUSTER"
python3 step6_excel_builder.py --cluster "$CLUSTER"

echo "🏁 === [Analysis Pipeline 완수] 격리 분석 가공 원부 배포 성료. ==="

0개의 댓글