【发布时间】: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