【问题标题】:spark-redshift takes a lot of time to write to redshiftspark-redshift 需要很长时间才能写入 redshift
【发布时间】:2016-06-14 20:45:58
【问题描述】:

我正在使用 kinesis 和 redshift 设置火花流光。我每 10 秒后从 kinesis 读取数据,对其进行处理并使用 spark-redshift lib 将其写入 redshift。

问题是只写 300 行要花很多时间。

这就是它在控制台中显示的内容

[Stage 56:====================================================> (193 + 1) / 200]

查看我的日志 df.write.format 正在这样做。

我在一台具有 4 gb ram 和 2 个核心 amazon EC2 的机器上设置了 spark 设置,以 --master local[*] 模式运行。

这是我创建流的方式

kinesisStream = KinesisUtils.createStream(ssc, APPLICATION_NAME, STREAM_NAME, ENDPOINT, REGION_NAME, INITIAL_POS, CHECKPOINT_INTERVAL, awsAccessKeyId =AWSACCESSID, awsSecretKey=AWSSECRETKEY, storageLevel=STORAGE_LEVEL)    
CHECKPOINT_INTERVAL = 60
storageLevel = memory

kinesisStream.foreachRDD(writeTotable)
def WriteToTable(df, type):
    if type in REDSHIFT_PAGEVIEW_TBL:
        df = df.groupby([COL_STARTTIME, COL_ENDTIME, COL_CUSTOMERID, COL_PROJECTID, COL_FONTTYPE, COL_DOMAINNAME, COL_USERAGENT]).count()
        df = df.withColumnRenamed('count', COL_PAGEVIEWCOUNT)

        # Write back to a table

        url = ("jdbc:redshift://" + REDSHIFT_HOSTNAME + ":" + REDSHIFT_PORT + "/" +   REDSHIFT_DATABASE + "?user=" + REDSHIFT_USERNAME + "&password="+ REDSHIFT_PASSWORD)

        s3Dir = 's3n://' + AWSACCESSID + ':' + AWSSECRETKEY + '@' + BUCKET + '/' + FOLDER

        print 'Start writing to redshift'
        df.write.format("com.databricks.spark.redshift").option("url", url).option("dbtable", REDSHIFT_PAGEVIEW_TBL).option('tempdir', s3Dir).mode('Append').save()

        print 'Finished writing to redshift'

请告诉我花这么多时间的原因

【问题讨论】:

    标签: apache-spark spark-streaming amazon-redshift


    【解决方案1】:

    在通过 Spark 和直接向 Redshift 写信时,我也有过类似的经历。 spark-redshift 将始终将数据写入 S3,然后使用 Redshift 复制功能将数据写入目标表。这种方法是写入大量记录的最佳实践和最有效的方法。这种方法也会给写入带来很多开销,尤其是当每次写入的记录数相对较少时。

    查看上面的输出,您似乎有大量的分区(可能有 200 个左右)。这可能是因为 spark.sql.shuffle.partitions 设置默认设置为 200。您可以找到更多详细信息in the Spark documentation

    组操作可能会生成 200 个分区。这意味着您正在对 S3 执行 200 次单独的复制操作,每个操作在获取连接和完成写入方面都有相当大的相关延迟。

    正如我们在下面的 cmets 和聊天中讨论的那样,您可以将分组结果合并到更少的分区中,对上面的行进行以下更改:

    df = df.coalesce(4).withColumnRenamed('count', COL_PAGEVIEWCOUNT)
    

    这会将分区数量从 200 个减少到 4 个,并将副本到 S3 的开销减少几个数量级。您可以试验分区数以优化性能。您还可以更改 spark.sql.shuffle.partitions 设置,以根据您正在处理的数据大小和可用内核数量来减少分区数量。

    【讨论】:

    • 不要你只写 3 行需要大约 4 分钟的时间。此外,即使我有 5000 行要写,4 分钟仍然是很多时间
    • 哇,我没有意识到需要这么长时间。在这种情况下,可能发生的情况是您有太多分区(从上面的输出中似乎就是这种情况)。这可能会导致从您的机器写入 S3 的瓶颈。我不确定这是否适用于流媒体,但对于常规的 spark 作业,例如 df.coalesce(1).write.format("com.databricks.spark.redshift").option("url", url)。 option("dbtable", REDSHIFT_PAGEVIEW_TBL).option('tempdir', s3Dir).mode('Append').save() 会起作用。您可以使用要合并的分区数量。
    • 我试过了,使用了 coalesce(4) 和缓存,但花了同样的时间。这很奇怪,但 4 分钟就像写 10 条记录或 1000 条记录一样多。我尝试联系 AWS,但也无济于事。尝试使用命令将 csv 直接从 s3 加载到 redshift 以查看是否需要时间,但这也需要几秒钟。
    • 你能帮帮我吗,我该如何调试它。在某处写日志以找出需要更多时间的地方。罪魁祸首肯定是 df.save,但在哪里?
    • 如果你把它贴在某个地方,我会看一下日志输出。这确实很奇怪。
    【解决方案2】:

    你是databrick API吗?这是已知问题。我有同样的问题。我确实与 Databric API 团队交谈过。从 Avaro 文件加载时,redshift 似乎没有提供良好的性能。我们确实与 AWS 团队进行了交谈。他们正在努力。 Databrick API 正在 S3 上创建 avaro 文件,然后复制命令将加载 avaro 文件。那是性能杀手。

    【讨论】:

    • 请将此作为评论发表
    猜你喜欢
    • 1970-01-01
    • 2013-10-11
    • 2012-12-04
    • 2012-12-03
    • 2014-07-21
    • 2019-10-26
    • 2012-06-24
    • 2013-11-15
    • 2014-02-23
    相关资源
    最近更新 更多