【问题标题】:Bad performance using spark against cassandra使用 spark 对 cassandra 的性能不佳
【发布时间】:2015-11-15 07:20:18
【问题描述】:

目前,我们正在挑战我们的架构,同时对我们的 cassandra 数据库使用 apache spark,因为我们遇到了非常糟糕的读取性能。

spark 和 cassandra 发生的硬件是一个云服务器,具有 16GB 内存和 8 个内核,并且使用 SSD 作为操作系统。

Cassandra 'data_file_directories' 设置为另一个硬盘,其测试结果为 hdparm -tT

Timing cached reads:   13140 MB in  1.99 seconds = 6604.42 MB/sec
Timing buffered disk reads: 428 MB in  3.00 seconds = 142.65 MB/sec

cassandra中的目标cf:

CREATE TABLE test.stats (
day timestamp,
received timestamp,
target inet,
via inet,
prefix blob,
rtt decimal,
PRIMARY KEY (day, received, target, via)
) WITH CLUSTERING ORDER BY (received ASC, target ASC, via ASC)
    AND bloom_filter_fp_chance = 0.01
    AND caching = '{"keys":"NONE", "rows_per_partition":"NONE"}'
    AND comment = ''
    AND compaction = {'min_threshold': '4', 'class':     'org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy',     'max_threshold': '32'}
    AND compression = {}
    AND dclocal_read_repair_chance = 0.1
    AND default_time_to_live = 0
    AND gc_grace_seconds = 864000
    AND max_index_interval = 2048
    AND memtable_flush_period_in_ms = 0
    AND min_index_interval = 128
    AND read_repair_chance = 0.0
    AND speculative_retry = '99.0PERCENTILE';

我们目前正在使用带有 java datastax 驱动程序 (2.1.6) 和 spark java 连接器 (spark-cassandra-connector_2.10,版本 1.4.0-M2) 的 Cassandra 2.1.6。

Spark 进程当前有一个设置了conf.set("spark.executor.memory", "2G"); 的工作人员。

通过提交/或序列化驱动程序启动一个简单的 spark 作业以读取一个显式分区键的所有行(大约有 83.520.000 行)需要 17 分钟。

该作业只是将所有行写入一个文件,该文件的最终大小为 1.2G。

这里是驱动代码:

CassandraTableScanJavaRDD<EchoRepliesBean> cassandraTable2 = null;
    switch (timespanMode)
    {
        case SIX_HOURS:
            Calendar calendarDay = Calendar.getInstance(TimeZone.getTimeZone("UTC"));
            calendarDay.setTimeInMillis(now);
            calendarDay.set(Calendar.HOUR_OF_DAY, 0);
            calendarDay.set(Calendar.MINUTE, 0);
            calendarDay.set(Calendar.SECOND, 0);
            calendarDay.set(Calendar.MILLISECOND, 0);
            Timestamp tsEnd = new Timestamp(calendarDay.getTimeInMillis());
            calendarDay.add(Calendar.DAY_OF_MONTH, -1);
            Timestamp tsStart = new Timestamp(calendarDay.getTimeInMillis());
            System.out.println(tsStart);
            cassandraTable2 = javaFunctions(sc).cassandraTable("test", "stats", mapRowTo(EchoRepliesBean.class))
                                               .where("day = ?", tsStart);
            break;
        default:
            /* make compiler happy */
            // cassandraTable = null;
    } 
    cassandraTable2.saveAsTextFile("/opt/out_TEST_" + System.currentTimeMillis());

    sc.stop();

这很奇怪,任何有关进一步调试的帮助或想法将不胜感激。

【问题讨论】:

  • 您是否尝试过使用较小的分区?分区有一定行数后性能是否突然下降?火花工作者是否在 Cassandra 节点上运行?您是否尝试增加分配给 spark worker 的内存?

标签: java cassandra apache-spark


【解决方案1】:

提高性能的一些可能方法:

  1. 通过在多个节点上对数据进行分区来提高并行度。由于您按天进行分区,因此您在一个节点上的一个分区中有大量行。这迫使读取和写入成为串行操作。如果您按小时进行分区,那么您的数据可能会分布在多个节点和多个 spark worker 中。

  2. 我怀疑您的 day 分区太大而无法放入单个 spark worker 的内存中,这可能会导致一些数据交换到磁盘。使用更小的分区、为 spark worker 提供更多内存或使用更多 spark worker 可以避免这种情况。

  3. 确保您的 spark worker 在 Cassandra 节点上运行,而不是在单独的机器上。如果工作人员在不同的机器上,那么将数据从节点转移到工作人员会产生大量网络开销。

  4. 确保您的云服务器为 Cassandra 使用本地存储,而不是网络存储。

为了调试,我会尝试在只有一行的分区上运行您的测试。如果这表现不佳,那么您的机器设置有问题。如果效果良好,则增加分区中的行数,直到您看到性能急剧下降。

【讨论】:

  • 1.这将是我们将在中期做的事情。 2. 我们将分区更改为每小时,结果大约有 200 万行。这将持续时间缩短到 26 秒! :) 3. 我们已经有了。 4. 迪托。谢谢吉姆!这有很大帮助。仍然困扰我的是由此产生的“速度”。输出文件的大小约为 100MByte,花费 26 秒并不多。
  • 你觉得26s还是太慢了?请记住,spark 仍然是批处理作业,而不是实时查询。旋转作业会产生开销。
  • 问题是我们没有经验,这就是我问的原因。 :)
【解决方案2】:

在 Cassandra 2.1.5 中引入了一个错误,并在 2.1.8 中进行了纠正,该错误导致火花性能显着下降:https://issues.apache.org/jira/browse/CASSANDRA-9637

升级到最新版本 (2.1.8) 可能会大大提高您的 spark 相关性能。

【讨论】:

  • 感谢 Jeff,我们更新了 cassandra。但我可能不会说这是否会产生重大影响,或者较小的分区是否会对性能产生更大的影响。
猜你喜欢
  • 2015-11-10
  • 2019-01-04
  • 2017-02-10
  • 2016-06-01
  • 2020-10-02
  • 2016-07-12
  • 2016-02-10
  • 1970-01-01
  • 2016-07-25
相关资源
最近更新 更多