【发布时间】:2020-02-16 16:15:31
【问题描述】:
我刚接触 spark 并开始使用 pyspark,我正在学习使用 pyspark 将数据从 kafka 推送到 hive。
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode
from pyspark.sql.functions import *
from pyspark.streaming.kafka import KafkaUtils
from os.path import abspath
warehouseLocation = abspath("spark-warehouse")
spark = SparkSession.builder.appName("sparkstreaming").getOrCreate()
df = spark.read.format("kafka").option("startingoffsets", "earliest").option("kafka.bootstrap.servers", "kafka-server1:66,kafka-server2:66").option("kafka.security.protocol", "SSL").option("kafka.ssl.keystore.location", "mykeystore.jks").option("kafka.ssl.keystore.password","mykeystorepassword").option("subscribe","json_stream").load().selectExpr("CAST(value AS STRING)")
json_schema = df.schema
df1 = df.select($"value").select(from_json,json_schema).alias("data").select("data.*")
上述方法不起作用,但是在提取数据后,我想将数据插入到 hive 表中。
由于我是全新的,寻求帮助。 提前表扬! :)
【问题讨论】:
标签: pyspark spark-streaming spark-streaming-kafka