"""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:
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()."
)
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()