使用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)
错误原因
- Schema字段名无法匹配嵌套结构:直接使用
host.id这种带点的字段名,Spark无法识别这是嵌套字典的路径,只会在顶层字典中查找名为host.id的键,自然找不到数据返回null。 - 字段类型定义错误:
tags是元组类型,应该用ArrayType(StringType)而非StringType。当把元组强行转成String时,Spark会输出对象的内存引用(如[Ljava.lang.Object;@xxx)。 - 数据结构不匹配:传入
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
相关产品推荐
相关产品推荐

