首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >如何使用PySpark向MySQL数据库流式传输数据?

如何使用PySpark向MySQL数据库流式传输数据?
EN

Stack Overflow用户
提问于 2018-11-13 02:50:11
回答 1查看 1.6K关注 0票数 1

我目前正在开发一个单页面web应用程序,它允许用户将大型CSV文件(目前正在测试一个~7 7GB的文件)上传到flask服务器,然后将该数据集流式传输到数据库。上传大约需要一分钟,文件将完全保存到flask服务器上的临时文件中。现在,我需要能够流式传输此文件并将其存储到数据库中。我做了一些研究,发现PySpark非常适合流式传输数据,我选择MySQL作为流式传输数据的数据库(但我对其他dbs和流式传输方法持开放态度)。我是一个初级开发人员,也是PySpark的新手,所以我不知道该怎么做。Spark streaming guide说数据必须通过Kafka、Flume、TCP socets等源获取,所以我想知道是否必须使用这些方法中的任何一种来将我的CSV文件导入到Spark中。然而,我遇到了这个great example,他们正在将csv数据流式传输到Azure SQL数据库中,看起来他们只是使用HDInsight直接读取文件,而不需要通过Kafka等流媒体源来摄取文件。唯一让我对这个例子感到困惑的是,他们正在使用Spark Spark集群将数据流式传输到数据库中,而我不确定如何将所有这些都整合到flask服务器中。我为缺少代码道歉,但目前我只有一个flask服务器文件,只有一条路径进行文件上传。任何示例、教程或建议都将不胜感激。

EN

回答 1

Stack Overflow用户

发布于 2018-11-13 19:18:40

我不确定流的部分,但spark可以有效地处理大文件-并且存储到数据库表将并行完成,所以在不太了解你的详细信息的情况下,如果你在服务器上有上传的文件,我会说:

如果我想在表中保存一个像csv这样的大型结构化文件,我会这样开始:

代码语言:javascript
复制
# start with some basic spark configuration, e.g. we want the timezone to be UTC 
conf = SparkConf()
conf.set('spark.sql.session.timeZone', 'UTC')
# this is important: you need to have the mysql connector jar for the right mysql version:
conf.set('jars', 'path to mysql connector jar you can download from here: https://dev.mysql.com/downloads/connector/odbc/')
# instantiate a spark session: the first time it will take a few seconds
spark = SparkSession.builder \
    .config(conf=conf) \
    .appName('Huge File uploader') \
    .getOrCreate()

# read the file first as a dataframe
df = spark.read.csv('path to 7GB/ huge csv file')

# optionally, add a filename column
from pyspark.sql import functions as F
df = df.withColumn('filename', F.lit('thecurrentfilename'))

# write it to the table
df.write.format('jdbc').options(
            url='e.g. localhost:port',
            driver='com.mysql.cj.jdbc.Driver',  # the driver for MySQL
            dbtable='the table name to save to',
            user='user',
            password='secret',
        ).mode('append').save()

注意这里的' append‘模式:这里的问题是spark不能在表上执行更新,它要么追加新行,要么替换表中的内容。

因此,如果您的csv是这样的:

代码语言:javascript
复制
id, name, address....

您最终将得到一个具有相同字段的表。

这是我能想到的最基本的例子,所以你可以从spark开始,而不考虑spark集群或任何其他相关的东西。我建议你带着这个去转转,看看这是否适合你的需求:)

此外,请记住,这可能需要几秒钟或更长时间,这取决于您的数据、数据库所在的位置、您的计算机和数据库加载,因此让api保持异步可能是一个好主意,同样,我不知道您的任何其他详细信息。

希望这能有所帮助。祝好运!

票数 1
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/53268342

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档