【发布时间】: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