【问题标题】: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()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-12-30
      • 2017-04-05
      • 2016-09-29
      • 2023-03-18
      • 2018-08-20
      • 1970-01-01
      • 2021-08-30
      • 1970-01-01
      相关资源
      最近更新 更多