如何用Python/Scala基于100k基准数据集生成500万相似唯一记录
Hey there! Let's tackle this problem—you want to scale your 100k unique dataset up to 5M similar records, with a specific column staying unique, no random garbage data. Perfect, here's how to do it in both Python and Scala/Spark, since both are great fits depending on your setup.
Python实现(用Pandas)
核心思路
We'll create controlled, small perturbations of your existing data to keep it similar, while using a unique identifier (either a new column or modifying an existing one) to ensure no duplicates in your specified column.
代码示例
This example uses a new unique_id column for uniqueness (adjust if you need to use an existing column like count):
import pandas as pd import numpy as np from uuid import uuid4 # 1. Load your original 100k dataset (adjust path/format as needed) df_original = pd.read_csv("original_data.csv") original_size = len(df_original) target_size = 5_000_000 need_generate = target_size - original_size # 2. Function to generate a similar row with tiny perturbations def generate_similar_row(row): # Add small random offset to lat/long (adjust range based on your data's precision) lat = row["latitude"] + np.random.uniform(-0.001, 0.001) lon = row["longitude"] + np.random.uniform(-0.001, 0.001) # Keep step/count same as original (or sample from original distribution if varied) step = row["step"] count = row["count"] # Generate a unique UUID for the identifier column unique_id = str(uuid4()) return pd.Series([lat, lon, step, count, unique_id], index=["latitude", "longitude", "step", "count", "unique_id"]) # 3. Batch-generate rows (vectorized approach for speed) # Sample original data repeatedly (with replacement) to get enough rows sample_batch = df_original.sample(n=need_generate, replace=True) # Apply perturbations in bulk sample_batch["latitude"] += np.random.uniform(-0.001, 0.001, size=need_generate) sample_batch["longitude"] += np.random.uniform(-0.001, 0.001, size=need_generate) # Add unique IDs sample_batch["unique_id"] = [str(uuid4()) for _ in range(need_generate)] # 4. Clean up and merge with original data # Drop any accidental duplicates (UUIDs are almost never duplicate, but just in case) df_generated = sample_batch.drop_duplicates(subset=["unique_id"]) # Add unique IDs to original data too df_original["unique_id"] = [str(uuid4()) for _ in range(original_size)] # Combine everything df_final = pd.concat([df_original, df_generated], ignore_index=True) # 5. Verify and save print(f"Final dataset size: {len(df_final)}") print(f"Unique column has no duplicates: {df_final['unique_id'].nunique() == len(df_final)}") df_final.to_csv("benchmark_data.csv", index=False)
关键调整 Tips
- Using an existing column for uniqueness: If you need
countto be unique, replace the UUID logic with an incrementing sequence. For example:max_count = df_original["count"].max() df_generated["count"] = range(max_count + 1, max_count + 1 + len(df_generated)) - Tweak perturbation range: If your original lat/long are more precise (e.g., 6+ decimal places), use a smaller offset like
±0.0001to keep data more similar. - Speed up generation: The vectorized approach above is way faster than looping—stick with it for large datasets.
Scala/Spark实现(适合大数据量)
If you're dealing with 5M records, Spark's distributed processing is a better fit (no memory issues!). Here's how to do it:
核心思路
We'll sample your original dataset repeatedly, apply small perturbations, and use Spark's built-in functions to generate globally unique identifiers.
代码示例
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ object BenchmarkDataGenerator { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("BenchmarkDataGenerator") .master("local[*]") // Remove this for cluster deployment .getOrCreate() import spark.implicits._ // 1. Load original dataset (adjust format/path) val originalDF = spark.read .option("header", "true") .csv("original_data.csv") .withColumn("latitude", $"latitude".cast(DoubleType)) .withColumn("longitude", $"longitude".cast(DoubleType)) .withColumn("step", $"step".cast(IntegerType)) .withColumn("count", $"count".cast(IntegerType)) val originalCount = originalDF.count().toInt val targetTotal = 5000000 val needGenerate = targetTotal - originalCount // 2. Generate similar data with perturbations val generatedDF = originalDF.sample(withReplacement = true, fraction = needGenerate.toDouble / originalCount) // Add ±0.001 offset to lat/long .withColumn("latitude", $"latitude" + (rand() * 0.002 - 0.001)) .withColumn("longitude", $"longitude" + (rand() * 0.002 - 0.001)) // Generate UUID for unique column .withColumn("unique_id", expr("uuid()")) // 3. Merge with original data (add unique IDs to original) val originalWithUnique = originalDF.withColumn("unique_id", expr("uuid()")) var finalDF = originalWithUnique.union(generatedDF).dropDuplicates("unique_id") // 4. Handle edge case: if deduplication reduced size, generate more val currentCount = finalDF.count().toInt if (currentCount < targetTotal) { val supplementDF = originalDF.sample(withReplacement = true, fraction = (targetTotal - currentCount).toDouble / originalCount) .withColumn("latitude", $"latitude" + (rand() * 0.002 - 0.001)) .withColumn("longitude", $"longitude" + (rand() * 0.002 - 0.001)) .withColumn("unique_id", expr("uuid()")) finalDF = finalDF.union(supplementDF).dropDuplicates("unique_id") } // 5. Save and verify finalDF.write.mode("overwrite").option("header", "true").csv("benchmark_data.csv") println(s"Final dataset size: ${finalDF.count()}") println(s"Unique column is duplicate-free: ${finalDF.select("unique_id").distinct().count() == finalDF.count()}") spark.stop() } }
关键调整 Tips
- Using existing column for uniqueness: If you need
countto be unique, replace the UUID with an incrementing sequence:val maxCount = originalDF.agg(max("count")).head().getInt(0) val generatedDFWithUniqueCount = generatedDF.withColumn("count", maxCount + row_number().over(Window.orderBy("unique_id"))) - Cluster deployment: Remove the
.master("local[*]")line when running on a Spark cluster—Spark will handle distributed processing automatically.
内容的提问来源于stack exchange,提问作者namrutha

