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

PySpark读取Kafka JSON转Integer列出现null值的解决方法

PySpark 转换ID列为整型返回null问题解决

问题场景

构建PySpark Structured Streaming作业消费Kafka的Twitter数据时,需要将tweet_id、userID两列修改为Integer类型。
初始实现代码如下:

import findspark
from pyspark import SparkConf, SparkContext
import pyspark
from pyspark.streaming import StreamingContext
from pyspark.sql.functions import from_json, col
from pyspark.sql import SparkSession
from pyspark.sql.types import *
import os

os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 pyspark-shell'
spark = SparkSession.builder.master("local[2]").appName('LearningDataframeWork').getOrCreate()

schema = StructType([ 
        StructField("tweet_id", StringType(), True),
        StructField("tweet_text" , StringType(), True),
        StructField("userID" , StringType(), True),
        StructField("username" , StringType(), True),
        ])


df_tweets = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "TwitterTweets") \
    .option("startingOffsets", "earliest") \
    .load() \
    .select(from_json(col("value").cast("string"),schema).alias("converted")) \
    .select(col("converted.*"))

df = df_tweets.select('*')


query = df.writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

query.awaitTermination()

初始运行可以正常输出字符串类型的ID数据,结果如下:

|           tweet_id|          tweet_text|             userID|      username|
+-------------------+--------------------+-------------------+--------------+
|1545020704969695236|RT @dme_363: Wait...|1341441279905902593|      dme__363|
|1545020704927625216|RT @DeadlineDayLi...|          634939603|    mattlobooo|
|1545020703350685698|@10Junioor Ronald...|          126401539|       IGAOC12|
|1545020702000128003|RT @ManUtdMEN: Po...|1247589115010351109|       D_Nyeko|
|1545020700334981121|RT @AgoluaVictor:...| 757194951985856512| Iam_Bwoiralph|
|1545020699764576256|Pré-saison avec M...|         2255497760|     sunusport|
|1545020696283209728|RT @ManuelMenacho...|         2305097051|SoetanAdebare1|
|1545020695477997570|@Ayden__x Ronaldo...|         4121380654|      blaq_gem|
|1545020691220766720|RT @JamesSunday_:...| 757194951985856512| Iam_Bwoiralph|
|1545020689522073601|RT @CFCBlues_com:...|1542264394675109888|      CBrayzay|

从Kafka消费到的原始JSON数据样例如下:

{"tweet_id": 15450003263664388, "tweet_text": "RT @Lordd: Who wears the Number 7 better??\\nLike for Kante, Retweet for Ronaldo\\n\\n||Ronaldo to Chelsea", "userID": 1196913590, "username": "davo_matsa"}

尝试使用.withColumn()方法将列转为Integer类型时,类型转换操作执行成功,但所有ID字段的值全部变为null,尝试代码如下:

df= df.withColumn("tweet_id_coverted", col("tweet_id").cast("Integer"))

转换后的输出结果:

+-------------------+--------------------+-------------------+--------------+-----------------+
|           tweet_id|          tweet_text|             userID|      username|tweet_id_coverted|
+-------------------+--------------------+-------------------+--------------+-----------------+
|1545020704969695236|RT @dme_363: Wait...|1341441279905902593|      dme__363|             null|
|1545020704927625216|RT @DeadlineDayLi...|          634939603|    mattlobooo|             null|
|1545020703350685698|@10Junioor Ronald...|          126401539|       IGAOC12|             null|
|1545020702000128003|RT @ManUtdMEN: Po...|1247589115010351109|       D_Nyeko|             null|
|1545020700334981121|RT @AgoluaVictor:...| 757194951985856512| Iam_Bwoiralph|             null|
|1545020699764576256|Pré-saison avec M...|         2255497760|     sunusport|             null|
|1545020696283209728|RT @ManuelMenacho...|         2305097051|SoetanAdebare1|             null|
|1545020695477997570|@Ayden__x Ronaldo...|         4121380654|      blaq_gem|             null|
|1545020691220766720|RT @JamesSunday_:...| 757194951985856512| Iam_Bwoiralph|             null|
|1545020689522073601|RT @CFCBlues_com:...|1542264394675109888|      CBrayzay|             null|

问题原因

Integer(32位整型)的取值范围为-2147483648 ~ 2147483647,即最大只能存储长度为10位的整数。样例中的tweet_id普遍为19位长度,数值远大于Integer类型的上限,数值溢出导致cast转换时直接返回null。

解决方案

  • Twitter的tweet_id、长用户ID均为64位整数,需要使用*LongType(长整型)*存储,Long类型取值范围为-9223372036854775808 ~ 9223372036854775807,完全可以覆盖这类长ID的存储需求。
  • 推荐直接在定义解析Schema时就将两个ID字段的类型指定为LongType,从源头按正确类型解析数据,无需后续二次转换:
schema = StructType([ 
        StructField("tweet_id", LongType(), True),
        StructField("tweet_text" , StringType(), True),
        StructField("userID" , LongType(), True),
        StructField("username" , StringType(), True),
        ])
  • 如果已经按String类型读取完数据,后续转换时将目标类型改为Long即可:
df = df.withColumn("tweet_id", col("tweet_id").cast("Long")) \
       .withColumn("userID", col("userID").cast("Long"))

注意:如果后续需要对接不支持64位整数的存储系统,可以将ID转为字符串类型存储,禁止使用Integer类型存储这类长ID,必然出现溢出问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 18:57:21