【发布时间】:2018-06-28 13:49:10
【问题描述】:
使用 Pyspark,如果数据帧 A 中的 IP 地址在 IP 网络范围内或命中数据帧 B 中的相同 IP 地址,我想加入/合并。
数据帧 A 仅包含 IP 地址,而另一个数据帧具有 IP 地址或带有 CIDR 的 IP 地址。这是一个例子。
Dataframe A
+---------------+
| ip_address|
+---------------+
| 192.0.2.2|
| 164.42.155.5|
| 52.95.245.0|
| 66.42.224.235|
| ...|
+---------------+
Dataframe B
+---------------+
| ip_address|
+---------------+
| 123.122.213.34|
| 41.32.241.2|
| 66.42.224.235|
| 192.0.2.0/23|
| ...|
+---------------+
那么预期的输出如下所示
+---------------+--------+
| ip_address| is_in_b|
+---------------+--------+
| 192.0.2.2| true| -> This is in the same network range as 192.0.2.0/23
| 164.42.155.5| false|
| 52.95.245.0| false|
| 66.42.224.235| true| -> This is in B
| ...| ...|
+---------------+--------+
我首先想尝试的想法是使用 udf 逐一比较并在 CIDR 出现时检查 IP 范围,但似乎 udf 没有多个数据帧。我还尝试将 df B 转换为列表,然后进行比较。但是,由于A行数*B行数超过1亿,效率非常低,耗时较长。有没有有效的解决方案?
编辑: 有关更多详细信息,我使用以下代码在不使用 pyspark 并使用任何库的情况下进行检查。
def cidr_to_netmask(c):
cidr = int(c)
mask = (0xffffffff >> (32 - cidr)) << (32 - cidr)
return (str((0xff000000 & mask) >> 24) + '.' + str((0x00ff0000 & mask) >> 16) + '.' + str((0x0000ff00 & mask) >> 8) + '.' + str((0x000000ff & mask)))
def ip_to_numeric(ip):
ip_num = 0
for i, octet in enumerate(ip.split('.')):
ip_num += int(octet) << (24 - (8 * i))
return ip_num
def is_in_ip_network(ip, network_addr):
if len(network_addr.split('/')) < 2:
return ip == network_addr.split('/')[0]
else:
network_ip, cidr = network_addr.split('/')
subnet = cidr_to_netmask(cidr)
return (ip_to_numeric(ip) & ip_to_numeric(subnet)) == (ip_to_numeric(network_ip) & ip_to_numeric(subnet))
【问题讨论】:
-
我不熟悉“在 IP 网络范围内”的逻辑——你能举例说明一下吗?还有为什么
52.95.245.0和66.42.224.235在输出中显示为false?这些显然在 B 中。我错过了什么吗? -
@pault 我只是想展示一个输出示例,所以我修改了示例数据帧。对于网络范围,gist.github.com/tott/7684443 即使在 PHP 中也是一个很好的例子。在 Python 中,我通常使用
netaddr库来执行与IPAddress(x) in IPNetwork(y)相同的操作。 netaddr.readthedocs.io/en/latest/tutorial_01.html
标签: python pyspark apache-spark-sql