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

如何基于StartIp/EndIp列高效扩展IP范围至Spark DataFrame?

高效展开Spark DataFrame中的IP范围(20亿级数据)

问题说明

现有一张包含StartIp、EndIp、value列的超大型数据表(20亿条记录),需要将每条记录中的IP范围展开为单个IP+对应value的行结构。

输入示例

StartIp,EndIp,value
1.0.0.9,1.0.0.18,a
172.0.0.2,180.0.0.1,c

期望输出示例

1.0.0.9,a
1.0.0.10,a
...
1.0.0.18,a
172.0.0.2,c
...
180.0.0.1,c

现有基础函数

已实现生成IP范围内所有IP的Scala工具函数:

import java.net.InetAddress
import org.apache.spark.sql.{Row, SparkSession}

object ExpandIpRange {
  def main(args: Array[String]): Unit = {
    // 示例IP范围
    val startIp = "192.168.0.1"
    val endIp = "192.168.0.5"

    // 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("ExpandIpRange")
      .getOrCreate()

    try {
      // 展开IP范围为序列
      def expandIpRange(startIp: String, endIp: String): Seq[String] = {
        val startAddress = InetAddress.getByName(startIp)
        val endAddress = InetAddress.getByName(endIp)
        val startInt = ipToInt(startAddress)
        val endInt = ipToInt(endAddress)
        (startInt to endInt).map(intToIp)
      }

      // IP转整数
      def ipToInt(address: InetAddress): Int = {
        val bytes = address.getAddress
        ((bytes(0) & 0xFF) << 24) |
          ((bytes(1) & 0xFF) << 16) |
          ((bytes(2) & 0xFF) << 8) |
          (bytes(3) & 0xFF)
      }

      // 整数转IP
      def intToIp(address: Int): String = {
        InetAddress.getByAddress(intToBytes(address)).getHostAddress
      }

      // 整数转字节数组
      def intToBytes(address: Int): Array[Byte] = {
        Array(
          ((address >> 24) & 0xFF).toByte,
          ((address >> 16) & 0xFF).toByte,
          ((address >> 8) & 0xFF).toByte,
          (address & 0xFF).toByte
        )
      }

      // 展开示例IP范围
      val ips = expandIpRange(startIp, endIp)
      val rows = ips.map(ip => Row(ip))
    } finally {
      spark.stop()
    }
  }
}

高效适配方案(针对20亿级数据)

直接用普通UDF展开会导致Executor内存溢出,必须采用分布式低内存的实现方式:

1. 先将IP转换为整数,计算范围边界

先把IP转为整数类型,避免在内存中生成全量IP序列:

import org.apache.spark.sql.functions._
import spark.implicits._

// 注册IP转整数的UDF
val ipToIntUdf = udf((ip: String) => {
  val address = InetAddress.getByName(ip)
  val bytes = address.getAddress
  ((bytes(0) & 0xFF) << 24) |
    ((bytes(1) & 0xFF) << 16) |
    ((bytes(2) & 0xFF) << 8) |
    (bytes(3) & 0xFF)
})

// 转换原数据表,添加整数形式的IP边界和范围长度
val ipRangeDF = originalDF
  .withColumn("startInt", ipToIntUdf(col("StartIp")))
  .withColumn("endInt", ipToIntUdf(col("EndIp")))

2. 用Spark内置函数分布式展开IP范围

利用sequence生成整数序列,结合explode实现分布式展开,避免单节点内存压力:

// 注册整数转IP的UDF
val intToIpUdf = udf((intIp: Int) => {
  val bytes = Array(
    ((intIp >> 24) & 0xFF).toByte,
    ((intIp >> 16) & 0xFF).toByte,
    ((intIp >> 8) & 0xFF).toByte,
    (intIp & 0xFF).toByte
  )
  InetAddress.getByAddress(bytes).getHostAddress
})

// 生成IP序列并展开为目标结构
val expandedDF = ipRangeDF
  .withColumn("ipInt", explode(sequence(col("startInt"), col("endInt"))))
  .withColumn("Ip", intToIpUdf(col("ipInt")))
  .select("Ip", "value")

3. 关键性能优化点

  • 分区调优:根据集群资源调整数据表分区数,建议每个分区处理100万-500万条IP记录
  • 内存配置:给Executor分配足够内存(如--executor-memory 16G),避免OOM
  • 序列化优化:启用Kryo序列化减少数据传输开销
  • 合并小文件:展开后数据量极大,用coalesce或repartition合并分区,减少输出文件数量

内容的提问来源于stack exchange,提问作者CompEng

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 22:23:36