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

基于Zeppelin实现IP关联国家代码:表连接方法咨询

在Zeppelin中实现IP地址与国家代码表的关联

嘿,你已经写好了IP转长整型的核心函数,接下来我们一步步把两张表关联起来,拿到IP对应的国家代码:

第一步:完善IP转换函数并注册为UDF

你的ipToLong函数逻辑没问题,我们给它加个空值判断避免报错,同时注册成Spark SQL的UDF,这样不管用DataFrame API还是SQL都能方便调用:

// 定义IP转长整型函数,增加空值处理
def ipToLong(dottedIP: String): Long = {
  if (dottedIP == null || dottedIP.trim.isEmpty) 0L
  else {
    val addrArray: Array[String] = dottedIP.split("\\.")
    var num: Long = 0
    var i: Int = 0
    while (i < addrArray.length) {
      val power: Int = 3 - i
      num = num + ((addrArray(i).toInt % 256) * Math.pow(256, power)).toLong
      i += 1
    }
    num
  }
}

// 注册为Spark SQL可用的UDF
spark.udf.register("ip_to_long", ipToLong _)

第二步:加载并预处理两张表

处理你的审计表(mamta_audit.csv)

假设这张表包含原始IP地址列,我们先把它转成DataFrame,并且生成IP的长整型列:

// 加载审计表RDD
val rdd1 = sc.textFile("/user/mamta/mamta_audit/mamta_audit.csv")
// 跳过表头(如果你的CSV没有表头,直接去掉下面两行)
val headerLine = rdd1.first()
val cleanAuditRdd = rdd1.filter(_ != headerLine)

// 转成DataFrame,这里假设IP是每行的第一列,根据你的实际结构调整索引
val auditDF = cleanAuditRdd.map(line => {
  val columns = line.split(",") // 逗号分隔,如果是其他分隔符改成对应的,比如制表符用"\\t"
  columns(0) // 提取IP地址
}).toDF("ip_address")

// 添加IP长整型列
val auditWithLongIP = auditDF.withColumn("ip_long", callUDF("ip_to_long", col("ip_address")))

加载IP段-国家代码映射表

假设你的另一张表是IP段和国家代码的映射(比如包含起始IP长、结束IP长、国家代码),同样转成DataFrame:

// 替换成你的映射表实际路径
val ipMappingRdd = sc.textFile("/user/mamta/ip_country_mapping.csv")
// 同样跳过表头(无表头则省略)
val mappingHeader = ipMappingRdd.first()
val cleanMappingRdd = ipMappingRdd.filter(_ != mappingHeader)

val ipMappingDF = cleanMappingRdd.map(line => {
  val parts = line.split(",")
  // 对应起始IP长、结束IP长、国家代码,根据你的表结构调整顺序
  (parts(0).toLong, parts(1).toLong, parts(2))
}).toDF("start_ip_long", "end_ip_long", "country_code")

第三步:关联两张表获取结果

我们用左连接来确保所有审计记录都被保留,没有匹配到IP段的记录国家代码会显示null,你可以选DataFrame API或者SQL两种方式:

方式1:DataFrame API关联

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

val finalResult = auditWithLongIP.join(
  ipMappingDF,
  auditWithLongIP("ip_long").between(ipMappingDF("start_ip_long"), ipMappingDF("end_ip_long")),
  "left"
).select(
  auditWithLongIP("ip_address"), // 保留原始IP
  ipMappingDF("country_code")    // 匹配到的国家代码
)

// 查看前20条结果
finalResult.show()

方式2:Spark SQL查询

如果更习惯写SQL,可以把两个DataFrame注册成临时视图,然后写SQL语句:

// 注册临时视图
auditWithLongIP.createOrReplaceTempView("audit_table")
ipMappingDF.createOrReplaceTempView("ip_country_map")

// 执行关联查询
val finalResult = spark.sql("""
  SELECT 
    a.ip_address,
    m.country_code
  FROM audit_table a
  LEFT JOIN ip_country_map m
  ON a.ip_long BETWEEN m.start_ip_long AND m.end_ip_long
""")

// 查看结果
finalResult.show()

一些小提示

  • 要是你的CSV分隔符不是逗号,记得把split(",")改成对应的分隔符,比如制表符用split("\\t")。
  • 如果IP映射表数据量很大,建议给start_ip_long和end_ip_long列做分区或者建立索引,能大幅提升关联速度。
  • 要是有格式错误的IP地址,可以在ipToLong里加try-catch块捕获异常,避免整个任务失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:16:35