You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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 count to 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.0001 to 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 count to 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 08:42:00