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

使用pyspark.createDataFrame转换ES来源RDD为DataFrame时返回空值问题

ElasticSearch RDD转DataFrame空值与类型错误问题解决

问题背景

从ElasticSearch获取的RDD结构如下(ES列表被转为元组):

[
    ('rty456ui', {'@timestamp': '2022-10-10T24:56:10.000259+0000', 'host': {'id': 'test-host-id-1'}, 'watchlists': {'ioc': {'summary': '127.0.0.1', 'tags': ('Dummy Tag',)}}, 'source': {'ip': '127.0.0.1'}, 'event': {'created': '2022-10-10T13:56:10+00:00', 'id': 'rty456ui'}, 'tags': ('Mon',)}),
    ('cxs980qw', {'@timestamp': '2022-10-10T13:56:10.000259+0000', 'host': {'id': 'test-host-id-2'}, 'watchlists': {'ioc': {'summary': '0.0.0.1', 'tags': ('Dummy Tag',)}}, 'source': {'ip': '0.0.0.1'}, 'event': {'created': '2022-10-10T24:56:10+00:00', 'id': 'cxs980qw'}, 'tags': ('Mon', 'Tue')})
]

期望转换为扁平结构的DataFrame:

+---------------+-----------+-----------+---------------------------+-----------------------+-----------------------+---------------+
|host.id        |event.id   |source.ip  |event.created              |watchlists.ioc.summary |watchlists.ioc.tags    |tags           |
+---------------+-----------+-----------+---------------------------+-----------------------+-----------------------+---------------+
|test-host-id-1 |rty456ui   |127.0.0.1  |2022-10-10T13:56:10+00:00  |127.0.0.1              |[Dummy Tag]            |[Mon]          |
|test-host-id-2 |cxs980qw   |0.0.0.1    |2022-10-10T24:56:10+00:00  |0.0.0.1                |[Dummy Tag]            |[Mon, Tue]     |
+---------------+-----------+-----------+---------------------------+-----------------------+-----------------------+---------------+

但实际得到的结果全为null,且tags字段显示对象引用:

+-------+--------+---------+-------------+----------------------+-------------------+-------------------------------+
|host.id|event.id|source.ip|event.created|watchlists.ioc.summary|watchlists.ioc.tags|tags                           |
+-------+--------+---------+-------------+----------------------+-------------------+-------------------------------+
|null   |null    |null     |null         |null                  |null               |[Ljava.lang.Object;@6c704e6e   |
|null   |null    |null     |null         |null                  |null               |[Ljava.lang.Object;@701ea4c8   |
+-------+--------+---------+-------------+----------------------+-------------------+-------------------------------+

原代码:

from pyspark.sql.types import StructType, StructField, StringType

schema = StructType([ 
    StructField("host.id",StringType(), True), 
    StructField("event.id",StringType(), True), 
    StructField("source.ip",StringType(), True), 
    StructField("event.created", StringType(), True), 
    StructField("watchlists.ioc.summary", StringType(), True), 
    StructField("watchlists.ioc.tags", StringType(), True),
    StructField("tags", StringType(), True)
  ])


df = spark.createDataFrame(es_rdd.map(lambda x: x[1]),schema)
df.show(truncate=False)

错误原因

  1. Schema字段名无法匹配嵌套结构:直接使用host.id这种带点的字段名,Spark无法识别这是嵌套字典的路径,只会在顶层字典中查找名为host.id的键,自然找不到数据返回null。
  2. 字段类型定义错误:tags是元组类型,应该用ArrayType(StringType)而非StringType。当把元组强行转成String时,Spark会输出对象的内存引用(如[Ljava.lang.Object;@xxx)。
  3. 数据结构不匹配:传入createDataFrame的是嵌套字典,但Schema定义的是扁平结构,两者无法自动映射。

解决方案

方法一:先扁平化RDD再生成DataFrame

手动提取嵌套字段,将元组转为列表,再用匹配的Schema创建DataFrame:

from pyspark.sql.types import StructType, StructField, StringType, ArrayType

# 扁平化每个元素,提取目标字段并转换元组为列表
flattened_rdd = es_rdd.map(lambda x: {
    "host.id": x[1]["host"]["id"],
    "event.id": x[1]["event"]["id"],
    "source.ip": x[1]["source"]["ip"],
    "event.created": x[1]["event"]["created"],
    "watchlists.ioc.summary": x[1]["watchlists"]["ioc"]["summary"],
    "watchlists.ioc.tags": list(x[1]["watchlists"]["ioc"]["tags"]),
    "tags": list(x[1]["tags"])
})

# 定义正确的Schema,数组类型对应列表/元组
schema = StructType([ 
    StructField("host.id", StringType(), True), 
    StructField("event.id", StringType(), True), 
    StructField("source.ip", StringType(), True), 
    StructField("event.created", StringType(), True), 
    StructField("watchlists.ioc.summary", StringType(), True), 
    StructField("watchlists.ioc.tags", ArrayType(StringType()), True),
    StructField("tags", ArrayType(StringType()), True)
])

df = spark.createDataFrame(flattened_rdd, schema)
df.show(truncate=False)

方法二:先创建嵌套结构DataFrame再展平

先定义匹配原始嵌套结构的Schema,再通过selectExpr提取目标字段:

from pyspark.sql.types import StructType, StructField, StringType, ArrayType

# 定义嵌套Schema,匹配原始字典结构
nested_schema = StructType([
    StructField("@timestamp", StringType(), True),
    StructField("host", StructType([
        StructField("id", StringType(), True)
    ]), True),
    StructField("watchlists", StructType([
        StructField("ioc", StructType([
            StructField("summary", StringType(), True),
            StructField("tags", ArrayType(StringType()), True)
        ]), True)
    ]), True),
    StructField("source", StructType([
        StructField("ip", StringType(), True)
    ]), True),
    StructField("event", StructType([
        StructField("created", StringType(), True),
        StructField("id", StringType(), True)
    ]), True),
    StructField("tags", ArrayType(StringType()), True)
])

# 创建嵌套结构的DataFrame
nested_df = spark.createDataFrame(es_rdd.map(lambda x: x[1]), nested_schema)

# 提取字段并重命名为扁平格式
flattened_df = nested_df.selectExpr(
    "host.id as `host.id`",
    "event.id as `event.id`",
    "source.ip as `source.ip`",
    "event.created as `event.created`",
    "watchlists.ioc.summary as `watchlists.ioc.summary`",
    "watchlists.ioc.tags as `watchlists.ioc.tags`",
    "tags"
)

flattened_df.show(truncate=False)

两种方法对比

  • 方法一:直接生成目标结构,代码简洁,适合明确知道需要提取的字段的场景。
  • 方法二:保留原始嵌套数据,后续可灵活提取其他字段,扩展性更强。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 07:25:17