26Z07f8a

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

import argparse
import sys
import pymysql


def run_reconciliation(args):
    print(
        f"[1/3] 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:
            # 1. 활성 객체 Upsert (Key, Size, ETag를 백틱으로 감쌈)
            print(
                f"[2/3] Upserting Active Objects directly from S3 Parquet ({s3_path})..."
            )

            insert_sql = f"""
                INSERT INTO {args.table_name} (
                    object_id, 
                    bucket, 
                    object_key, 
                    prefix, 
                    size_bytes, 
                    etag, 
                    event_time, 
                    __deleted
                )
                SELECT 
                    md5(concat(%s, ':', url_decode(`Key`))) AS object_id,
                    %s AS bucket,
                    url_decode(`Key`) AS object_key,
                    coalesce(regexp_extract(url_decode(`Key`), '^(.*)/[^/]+$', 1), '') AS prefix,
                    CAST(coalesce(`Size`, 0) AS BIGINT) AS size_bytes,
                    coalesce(`ETag`, '') AS etag,
                    %s AS event_time,
                    0 AS __deleted
                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.use_instance_profile" = "false"
                );
            """

            affected_rows = cursor.execute(
                insert_sql,
                (args.target_bucket, args.target_bucket, args.snapshot_time),
            )
            print(
                f" -> Successfully Upserted {affected_rows} active rows via FILES()."
            )

            # 2. Reconciliation (Missing Objects 백틱 적용 및 LEFT ANTI JOIN)
            print(
                "[3/3] 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 (
                    SELECT url_decode(`Key`) AS object_key
                    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.use_instance_profile" = "false"
                    )
                ) 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.")

    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개의 댓글