EC256/openenv-data-engineering
0
1"""2Generate advanced, multi-table corrupted datasets for the OpenEnv environment.3Includes schema drift, join corruption, nested JSON corruption, and misleading logs.4Uses Faker for deterministic, reproducible data.5"""6 7import os8import csv9import json10import random11from datetime import datetime, timedelta12 13from faker import Faker14 15SEED = 4216DATA_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "data")17NUM_ROWS = 50018 19os.makedirs(DATA_DIR, exist_ok=True)20 21 22def generate_users():23 """24 Users dataset with:25 - Schema drift: user_id as int vs string (30% drift)26 - Null corruption: 10% missing emails27 - Date format inconsistencies: 3 different formats28 """29 fake = Faker("en_US")30 random.seed(SEED)31 Faker.seed(SEED)32 33 users = []34 for i in range(NUM_ROWS):35 user_id = i if random.random() > 0.3 else f"USR_{i}"36 37 users.append(38 {39 "user_id": user_id,40 "name": fake.name(),41 "email": fake.email() if random.random() > 0.1 else None,42 "signup_date": fake.date_this_decade().strftime(43 random.choice(["%Y-%m-%d", "%d/%m/%Y", "%m-%d-%Y"])44 ),45 "status": random.choice(["active", "inactive", "suspended"]),46 }47 )48 49 fpath = os.path.join(DATA_DIR, "users.csv")50 with open(fpath, "w", newline="") as f:51 writer = csv.DictWriter(52 f, fieldnames=["user_id", "name", "email", "signup_date", "status"]53 )54 writer.writeheader()55 writer.writerows(users)56 print(f"[DATA] Generated {fpath} ({len(users)} rows)")57 return users58 59 60def generate_transactions(users):61 """62 Transactions dataset with:63 - Join corruption: 40% user_id mismatch (int vs USR_x format)64 - Nested JSON corruption: metadata always stored as JSON string65 - Null amounts: 10% missing66 - Silent join failure risk67 """68 fake = Faker("en_US")69 random.seed(SEED + 1)70 Faker.seed(SEED + 1)71 72 transactions = []73 for i in range(NUM_ROWS):74 user_ref = random.randint(0, NUM_ROWS - 1)75 76 if random.random() < 0.4:77 user_ref = f"USR_{user_ref}"78 79 amount = round(random.uniform(10, 500), 2)80 81 metadata = {82 "payment_method": random.choice(["card", "upi", "bank"]),83 "location": fake.city(),84 "device": random.choice(["mobile", "desktop", "tablet"]),85 }86 87 # Always serialize metadata as JSON string88 metadata = json.dumps(metadata)89 90 if random.random() < 0.1:91 amount = None92 93 transactions.append(94 {95 "transaction_id": i,96 "user_id": user_ref,97 "amount": amount,98 "metadata": metadata,99 "timestamp": (100 datetime(2024, 1, 1) + timedelta(minutes=random.randint(0, 525600))101 ).isoformat(),102 }103 )104 105 fpath = os.path.join(DATA_DIR, "transactions.csv")106 with open(fpath, "w", newline="") as f:107 writer = csv.DictWriter(108 f,109 fieldnames=["transaction_id", "user_id", "amount", "metadata", "timestamp"],110 )111 writer.writeheader()112 writer.writerows(transactions)113 print(f"[DATA] Generated {fpath} ({len(transactions)} rows)")114 return transactions115 116 117def generate_logs():118 """119 Logs dataset with:120 - Misleading noise (common false signals)121 - Rare true signal (schema drift warning, only 5 occurrences)122 - Tests agent ability to filter signal from noise123 """124 random.seed(SEED + 2)125 126 logs = []127 noise_messages = [128 "Null values detected in transactions.amount",129 "Data pipeline executed successfully",130 "Minor delay in ingestion",131 "Schema validation passed",132 "Cache refresh completed",133 "Rate limit threshold approaching",134 "Connection pool utilization at 60%",135 ]136 137 for _ in range(200):138 logs.append(139 {140 "level": random.choice(["INFO", "WARNING", "ERROR"]),141 "message": random.choice(noise_messages),142 "timestamp": (143 datetime(2024, 1, 1) + timedelta(minutes=random.randint(0, 525600))144 ).isoformat(),145 }146 )147 148 for _ in range(5):149 logs.append(150 {151 "level": "WARNING",152 "message": "User ID format mismatch detected between tables",153 "timestamp": (154 datetime(2024, 6, 1) + timedelta(minutes=random.randint(0, 525600))155 ).isoformat(),156 }157 )158 159 random.shuffle(logs)160 161 fpath = os.path.join(DATA_DIR, "logs.csv")162 with open(fpath, "w", newline="") as f:163 writer = csv.DictWriter(f, fieldnames=["level", "message", "timestamp"])164 writer.writeheader()165 writer.writerows(logs)166 print(f"[DATA] Generated {fpath} ({len(logs)} rows)")167 168 169if __name__ == "__main__":170 print("[START] Generating advanced datasets...")171 users = generate_users()172 generate_transactions(users)173 generate_logs()174 print("[END] All datasets generated.")175 