【问题标题】:how can I read real-time updated data from postgresql(or mysql) through spark streaming(or kafka)?如何通过火花流(或 kafka)从 postgresql(或 mysql)读取实时更新数据?
【发布时间】:2018-11-01 19:52:23
【问题描述】:
我从 postgresql 获取实时更新的数据,并希望将实时数据流水线到固定模型,以通过 spark 流或 kafka 预测客户。
请推荐任何运行良好的博客和确切代码,或任何您知道的信息/建议。postgresql/mysql 实时更新数据到 python/java 环境也可以!谢谢!
或者也许只是没有办法实现这一点?
【问题讨论】:
标签:
java
python
database
apache-kafka
spark-streaming
【解决方案1】:
这是我的回答,希望能帮到你。
我的 spark 版本是 2.2.0。编程语言是 Python。
数据流从kafka到mysql,kafka版本为0.9。
注意:
用mysql和kafka一定要找到正确的jar,可以去官网找这个。
这样的代码:
from pyspark import SparkContext, Row
from pyspark.sql import SparkSession
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
# note: the mysql's driver is must be correct
def getSparkSessionInstance(sparkConf):
if ('sparkSessionSingletonInstance' not in globals()):
globals()['sparkSessionSingletonInstance'] = SparkSession\
.builder\
.config(conf=sparkConf)\
.getOrCreate()
return globals()['sparkSessionSingletonInstance']
if __name__ == "__main__":
# mysql config
url = "jdbc:mysql://your_server:3306/spark_test"
table_name = "word_info"
username = "root"
pasword = "root"
# spark context init
para_seconds = 10
sc = SparkContext(appName="PythonStreamingDirectKafkaWordCount")
ssc = StreamingContext(sc, para_seconds)
# receiver in kafka
brokers = 'kafka1:9092'
topic = 'two-two-para'
# get streaming datas from kafka
kvs = KafkaUtils.createDirectStream(ssc, [topic], {"metadata.broker.list": brokers})
lines = kvs.map(lambda x: x[1])
# Convert RDDs of the words DStream to DataFrame and run SQL query
def process(time, rdd):
print("========= %s =========" % str(time))
if (rdd.isEmpty()):
return
try:
# Get the singleton instance of SparkSession
spark = getSparkSessionInstance(rdd.context.getConf())
# Convert RDD[String] to RDD[Row] to DataFrame
rowRdd = rdd.map(lambda w: Row(word=w))
wordsDataFrame = spark.createDataFrame(rowRdd)
# Creates a temporary view using the DataFrame.
wordsDataFrame.createOrReplaceTempView("words")
# Do word count on table using SQL and print it
wordCountsDataFrame = \
spark.sql("select word, count(*) as word_count from words group by word")
wordCountsDataFrame.show()
wordCountsDataFrame.write \
.format("jdbc") \
.option("url", url) \
.option("driver", "org.mariadb.jdbc.Driver") \
.option("dbtable", table_name) \
.option("user", username) \
.option("password", pasword) \
.save(mode="append")
except Exception as e:
print("Some error happen!")
print(e)
lines.foreachRDD(process)
# start job
ssc.start()
ssc.awaitTermination()