【问题标题】:Spark JoinWithCassandraTable on TimeStamp partition key STUCK时间戳分区键 STUCK 上的 Spark JoinWithCassandraTable
【发布时间】:2016-01-24 14:13:47
【问题描述】:

我正在尝试使用以下方法过滤一个巨大的 C* 表的一小部分:

    val snapshotsFiltered = sc.parallelize(startDate to endDate).map(TableKey(_)).joinWithCassandraTable("listener","snapshots_tspark")

    println("Done Join")
    //*******
    //get only the snapshots and create rdd temp table
    val jsons = snapshotsFiltered.map(_._2.getString("snapshot"))
    val jsonSchemaRDD = sqlContext.jsonRDD(jsons)
    jsonSchemaRDD.registerTempTable("snapshots_json")

与:

    case class TableKey(created: Long) //(created, imei, when)--> created = partititon key | imei, when = clustering key

而 cassandra 表架构是:

CREATE TABLE listener.snapshots_tspark (
created timestamp,
imei text,
when timestamp,
snapshot text,
PRIMARY KEY (created, imei, when) ) WITH CLUSTERING ORDER BY (imei ASC, when ASC)
AND bloom_filter_fp_chance = 0.01
AND caching = '{"keys":"ALL", "rows_per_partition":"NONE"}'
AND comment = ''
AND compaction = {'min_threshold': '4', 'class': 'org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy', 'max_threshold': '32'}
AND compression = {'sstable_compression': 'org.apache.cassandra.io.compress.LZ4Compressor'}
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';

问题是 println 完成后进程冻结,在 spark master ui 上没有错误。

[Stage 0:>                                                                                                                                (0 + 2) / 2]

Join 不能使用时间戳作为分区键吗?为什么会结冰?

【问题讨论】:

  • 您是否检查过是否有足够的资源来运行作业?
  • @eliasah 是的。内存:总计 5.5 GB,已使用 512.0 MB
  • 如果集合 snapshotsFiltered 返回为空,下一阶段会卡住吗?
  • 不,这不是原因。它可能主要由于缺乏资源而卡住。这可能是由于您想要执行的查询计划的复杂性。
  • @eliasah 为什么它很复杂?它只假设在创建的时间戳上进行聚类,并将创建的时间戳 > startDate 和 startDate 带到 endDate = "+ startDate + " and created

标签: mysql scala cassandra apache-spark datastax-enterprise


【解决方案1】:

通过使用:

sc.parallelize(startDate to endDate)

将 startData 和 endDate 作为 Longs 从 Dates 生成的格式:

("yyyy-MM-dd HH:mm:ss")

我用 spark 构建了一个巨大的数组(100,000 多个对象)来加入 C* 表,它根本没有卡住 - C* 努力使连接发生并返回数据。

最后,我将范围更改为:

case class TableKey(created_dh: String)
val data = Array("2015-10-29 12:00:00", "2015-10-29 13:00:00", "2015-10-29 14:00:00", "2015-10-29 15:00:00")
val snapshotsFiltered = sc.parallelize(data, 2).map(TableKey(_)).joinWithCassandraTable("listener","snapshots_tnew")

现在没事了。

【讨论】:

    猜你喜欢
    • 2015-10-24
    • 1970-01-01
    • 2023-02-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-07-29
    相关资源
    最近更新 更多