26Z10a1

Young-Kyoo Kim·2026년 8월 10일
#!/usr/bin/env python3
"""MinIO AIStor Inventory Parquet -> StarRocks Reconciliation Pipeline (Staging Table Fix)"""

import argparse
import sys
import pymysql


def run_reconciliation(args):
    print(
        f"[1/4] Connecting to StarRocks MySQL Port ({args.sr_host}:{args.sr_mysql_port})..."
    )
    try:
        conn = pymysql.connect(
            host=args.sr_host,
            port=args.sr_mysql_port,
            user=args.sr_user,
            password=args.sr_password,
            db=args.db_name,
            autocommit=True,
            charset="utf8mb4",
        )
    except Exception as e:
        print(f"[ERROR] Failed to connect to StarRocks: {e}", file=sys.stderr)
        sys.exit(1)

    inv_key = args.inventory_key.rstrip("/")
    if not inv_key.endswith(".parquet"):
        s3_path = f"s3://{args.inventory_bucket}/{inv_key}/*.parquet"
    else:
        s3_path = f"s3://{args.inventory_bucket}/{inv_key}"

    aistor_endpoint = args.aistor_endpoint
    enable_ssl = "true" if args.aistor_secure else "false"

    try:
        with conn.cursor(pymysql.cursors.DictCursor) as cursor:
            # Step 1. S3 Parquet 데이터를 임시 테이블(tmp_aistor_inventory)로 1차 고속 로드
            print(
                f"[2/4] Loading S3 Parquet into StarRocks Temporary Staging Table ({s3_path})..."
            )

            # 세션 내 임시 테이블 cleanup
            cursor.execute("DROP TEMPORARY TABLE IF EXISTS tmp_aistor_inventory;")

            create_tmp_sql = f"""
                CREATE TEMPORARY TABLE tmp_aistor_inventory AS
                SELECT 
                    md5(concat(%s, ':', url_decode(`Key`))) AS object_id,
                    url_decode(`Key`) AS object_key,
                    coalesce(regexp_extract(url_decode(`Key`), '^(.*)/[^/]+$', 1), '') AS prefix,
                    CAST(coalesce(`Size`, 0) AS BIGINT) AS size_bytes
                FROM FILES(
                    "path" = "{s3_path}",
                    "format" = "parquet",
                    "aws.s3.endpoint" = "{aistor_endpoint}",
                    "aws.s3.access_key" = "{args.aistor_access_key}",
                    "aws.s3.secret_key" = "{args.aistor_secret_key}",
                    "aws.s3.enable_ssl" = "{enable_ssl}",
                    "aws.s3.enable_path_style_access" = "true",
                    "aws.s3.use_instance_profile" = "false"
                );
            """
            staged_rows = cursor.execute(create_tmp_sql, (args.target_bucket,))
            print(f" -> Staged {staged_rows} rows into temporary table.")

            # Step 2. 임시 테이블 -> 본 테이블로 활성 객체 Upsert
            print(
                f"[3/4] Upserting Active Objects from Staging Table into {args.table_name}..."
            )

            insert_sql = f"""
                INSERT INTO {args.table_name} (
                    object_id, 
                    bucket, 
                    object_key, 
                    prefix, 
                    size_bytes, 
                    etag, 
                    event_time, 
                    __deleted
                )
                SELECT 
                    object_id,
                    %s AS bucket,
                    object_key,
                    prefix,
                    size_bytes,
                    '' AS etag,
                    %s AS event_time,
                    0 AS __deleted
                FROM tmp_aistor_inventory;
            """
            affected_rows = cursor.execute(
                insert_sql, (args.target_bucket, args.snapshot_time)
            )
            print(f" -> Successfully Upserted {affected_rows} active rows.")

            # Step 3. Reconciliation (임시 테이블 대상 LEFT ANTI JOIN 수행)
            print(
                "[4/4] Detecting and marking deleted (missing) objects in StarRocks..."
            )

            delete_reconcile_sql = f"""
                INSERT INTO {args.table_name} (
                    object_id, 
                    bucket, 
                    object_key, 
                    prefix, 
                    size_bytes, 
                    etag, 
                    event_time, 
                    __deleted
                )
                SELECT 
                    db.object_id,
                    db.bucket,
                    db.object_key,
                    db.prefix,
                    0 AS size_bytes,
                    '' AS etag,
                    %s AS event_time,
                    1 AS __deleted
                FROM {args.table_name} db
                LEFT ANTI JOIN tmp_aistor_inventory inv 
                  ON db.object_key = inv.object_key
                WHERE db.bucket = %s
                  AND db.event_time <= %s
                  AND db.__deleted = 0;
            """

            deleted_cnt = cursor.execute(
                delete_reconcile_sql,
                (args.snapshot_time, args.target_bucket, args.snapshot_time),
            )
            print(f" -> Marked {deleted_cnt} missing objects as deleted.")

            # 임시 테이블 명시적 삭제
            cursor.execute("DROP TEMPORARY TABLE IF EXISTS tmp_aistor_inventory;")

    except Exception as e:
        print(f"[ERROR] Execution Failed: {e}", file=sys.stderr)
        sys.exit(1)
    finally:
        conn.close()

    print("\n=== Inventory Reconciliation Completed Successfully ===")


def main():
    parser = argparse.ArgumentParser(
        description="MinIO AIStor Inventory Parquet -> StarRocks Reconciliation Pipeline"
    )

    parser.add_argument(
        "--target-bucket",
        required=True,
        help="MinIO Bucket Name being reconciled",
    )
    parser.add_argument(
        "--inventory-bucket",
        required=True,
        help="AIStor Bucket containing the inventory Parquet file",
    )
    parser.add_argument(
        "--inventory-key",
        required=True,
        help="S3 Key/Prefix to the inventory Parquet file(s)",
    )
    parser.add_argument(
        "--snapshot-time",
        required=True,
        help="Inventory execution start timestamp (Format: 'YYYY-MM-DD HH:MM:SS')",
    )

    parser.add_argument(
        "--aistor-endpoint",
        default="http://aistor-s3-service:9000",
        help="AIStor S3 Endpoint accessible from StarRocks BE nodes",
    )
    parser.add_argument(
        "--aistor-access-key", default="minioadmin", help="AIStor Access Key"
    )
    parser.add_argument(
        "--aistor-secret-key", default="minioadmin", help="AIStor Secret Key"
    )
    parser.add_argument(
        "--aistor-secure",
        action="store_true",
        help="Use HTTPS for AIStor connection",
    )

    parser.add_argument(
        "--sr-host", default="localhost", help="StarRocks FE Host IP/Domain"
    )
    parser.add_argument(
        "--sr-mysql-port",
        type=int,
        default=9030,
        help="StarRocks MySQL Port (default: 9030)",
    )
    parser.add_argument(
        "--sr-user", default="root", help="StarRocks DB User"
    )
    parser.add_argument(
        "--sr-password", default="", help="StarRocks DB Password"
    )
    parser.add_argument(
        "--db-name", default="minio_metadata", help="StarRocks Database Name"
    )
    parser.add_argument(
        "--table-name",
        default="object_current_state",
        help="StarRocks Table Name",
    )

    args = parser.parse_args()
    run_reconciliation(args)


if __name__ == "__main__":
    main()

0개의 댓글