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
相关产品推荐
相关产品推荐

