【问题标题】:Pyspark split the spark dataframe of type stringPyspark 拆分字符串类型的 spark 数据帧
【发布时间】:2020-01-02 03:41:49
【问题描述】:

我正在通过使用 spark(批处理,而不是流式处理)从 kafka 主题读取数据来创建 spark 数据帧。我想使用 spark 将此数据帧加载到 cassandra。数据帧采用如下字符串格式。

根 |-- 值:字符串(可为空=真)

+--------------------+
|value               |
+--------------------+
|"1,Visa,6574"       |
|"3,Visa,6574"       |
|"4,MasterCard,6574" |
|"5,MasterCard,6574" |
|"8,Maestro,8372"    |
+--------------------+

我尝试使用“,”分隔符拆分数据帧记录并形成新的数据帧,我可以将其数据传输到 cassandra。

创建如下 sparkDF。

df = spark \
.read \
.format("kafka") \
.option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS_CONS) \
.option("subscribe", KAFKA_TOPIC_NAME_CONS) \
.option("startingOffsets", "earliest") \
.load()
df2=df.selectExpr("CAST(value AS STRING)")
df2.printSchema()

我尝试使用 ',' 分割数据。

split_col=split(df2['value'],',')
df3=df2.withColumn('Name1',split_col.getItem(0))
df3=df2.withColumn('Name2',split_col.getItem(1))
df3=df2.withColumn('Name3',split_col.getItem(2))

上面的代码没有给出预期的结果,我得到了喜欢

根 |-- 值:字符串(可为空=真) |-- Name3: 字符串 (nullable = true)

+-------------------+-----+
|value              |Name3|
+-------------------+-----+
|"1,Visa,6574"      |6574"|
|"3,Visa,6574"      |6574"|
|"4,MasterCard,6574"|6574"|
|"5,MasterCard,6574"|6574"|
|"8,Maestro,8372"   |8372"|
+-------------------+-----+

我想得到如下的 put:

+-------------------+----------+------+
|Name1              |Name2     |Name3 |
+-------------------+----------+------+
| 1                 |Visa      |6574  |
| 3                 |Visa      |6574  |
| 4                 |MasterCard|6574  |
| 5                 |MasterCard|6574  |
| 8                 |Maestro   |8372  |
+-------------------+----------+------+

请帮忙!!

【问题讨论】:

    标签: python apache-spark pyspark apache-kafka apache-spark-sql


    【解决方案1】:

    您的解决方案完全没问题。唯一的问题是df2df3 在进行拆分并用于下一步之后的分配。在进行第一次拆分后,您分配给df3,但对于后续拆分,您仅使用df2。因此,只有第 3 个 split 语句被 spark 评估。

    明智的解决方案是在最后一次拆分之前不要分配给新变量

    df3 = df2.withColumn('Name1', f.split('value', ',').getItem(0)).\
                     withColumn('Name2', f.split('value', ',').getItem(1)).\
                     withColumn('Name3', f.split('value', ',').getItem(2))
    
    df3.show()
    +-----------------+-----+----------+-----+
    |            value|Name1|     Name2|Name3|
    +-----------------+-----+----------+-----+
    |      1,Visa,6574|    1|      Visa| 6574|
    |      3,Visa,6574|    3|      Visa| 6574|
    |4,MasterCard,6574|    4|MasterCard| 6574|
    |5,MasterCard,6574|    5|MasterCard| 6574|
    |   8,Maestro,8372|    8|   Maestro| 8372|
    +-----------------+-----+----------+-----+
    
    

    或在下一次拆分中使用分配的变量(除非必要,否则不鼓励使用这种方式)

    df3 = df2.withColumn('Name1', f.split('value', ',').getItem(0))
    
    df3 = df3.withColumn('Name2', f.split('value', ',').getItem(1))
    
    df3 = df3.withColumn('Name3', f.split('value', ',').getItem(2))
    
    df3.show()
    +-----------------+-----+----------+-----+
    |            value|Name1|     Name2|Name3|
    +-----------------+-----+----------+-----+
    |      1,Visa,6574|    1|      Visa| 6574|
    |      3,Visa,6574|    3|      Visa| 6574|
    |4,MasterCard,6574|    4|MasterCard| 6574|
    |5,MasterCard,6574|    5|MasterCard| 6574|
    |   8,Maestro,8372|    8|   Maestro| 8372|
    +-----------------+-----+----------+-----+
    
    

    【讨论】:

      猜你喜欢
      • 2021-04-16
      • 1970-01-01
      • 1970-01-01
      • 2017-03-10
      • 1970-01-01
      • 2021-12-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多