【问题标题】:Spark connector: Partition usage and performance issueSpark 连接器:分区使用和性能问题
【发布时间】:2015-04-12 09:39:48
【问题描述】:

我正在尝试运行一个 spark 作业(与 Ca​​ssandra 对话)来读取数据,进行一些聚合,然后将聚合写入 Cassandra

  • 我有 2 个表(monthly_active_users (MAU) , daily_user_metric_aggregates (DUMA))
  • 对于 MAU 中的每条记录,DUMA 中都会有一条或多条记录
  • 获取 MAU 中的每条记录并获取其中的 user_id,然后在 DUMA 中为该用户查找记录(应用服务器端过滤器,如 ('ms', 'md') 中的 metric_name
  • 如果指定 where 子句的 DUMA 中有一条或多条记录,那么我需要增加 appMauAggregate 映射的计数(应用明智的 MAU 计数)
  • 我测试了这个算法,按预期工作,但我想知道

1)它是一种优化的算法(或)有没有更好的方法呢?我有一种不正确的感觉,我没有看到加速。看起来 Cassandra 客户端正在为每个火花动作(收集)创建和关闭。处理小数据集需要很长时间。

2) Spark worker 不与 cassandra 位于同一位置,这意味着 spark worker 运行在与 C* 节点不同的节点(容器)中(我们可以将 spark worker 移动到 C* 节点以获取数据局部性)

3) 我看到正在为每个 spark 操作(收集)创建/提交 spark 作业,我相信这是 spark 的预期行为,无论如何都要减少从 C* 读取并创建连接以便数据检索快吗?

4) 这个算法的缺点是什么?您能否推荐更好的设计方法,即 w/r/t 分区策略、将 C* 分区加载到 Spark 分区、执行程序/驱动程序的内存要求?

5) 只要算法和设计方法都很好,那么我就可以进行火花调整。我正在使用 5 名工人(每个工人有 16 个 CPU 和 64GB 内存)

C* 架构:

月活跃用户数:

CREATE TABLE analytics.monthly_active_users ( 
    month text, 
    app_id uuid,
    user_id uuid, 
    PRIMARY KEY (month, app_id, user_id) 
) WITH CLUSTERING ORDER BY (app_id ASC, user_id ASC)

数据:

cqlsh:analytics> select * from monthly_active_users limit 2;   
 month  | app_id                               | user_id 

--------+--------------------------------------+-------------------------------------- 
 2015-2 | 108eeeb3-7ff1-492c-9dcd-491b68492bf2 | 199c0a31-8e74-46d9-9b3c-04f67d58b4d1 
 2015-2 | 108eeeb3-7ff1-492c-9dcd-491b68492bf2 | 2c70a31a-031c-4dbf-8dbd-e2ce7bdc2bc7 

杜马:

CREATE TABLE analytics.daily_user_metric_aggregates ( 
    metric_date timestamp, 
    user_id uuid,
    metric_name text, 
    "count" counter, 
    PRIMARY KEY (metric_date, user_id, metric_name)
) WITH CLUSTERING ORDER BY (user_id ASC, metric_name ASC) 

数据:

cqlsh:analytics> select * from daily_user_metric_aggregates where metric_date='2015-02-08' and user_id=199c0a31-8e74-46d9-9b3c-04f67d58b4d1; 
 metric_date | user_id                                                         | metric_name       | count 
--------------------------+--------------------------------------+-------------------+------- 
 2015-02-08 | 199c0a31-8e74-46d9-9b3c-04f67d58b4d1 | md                      |     1     
 2015-02-08 | 199c0a31-8e74-46d9-9b3c-04f67d58b4d1 | ms                      |     1 

火花工作:

import java.net.InetAddress 
import java.util.concurrent.atomic.AtomicLong 
import java.util.{Date, UUID} 

import com.datastax.spark.connector.util.Logging 
import org.apache.spark.{SparkConf, SparkContext} 
import org.joda.time.{DateTime, DateTimeZone} 

import scala.collection.mutable.ListBuffer 

object MonthlyActiveUserAggregate extends App with Logging { 

    val KeySpace: String = "analytics" 
    val MauTable: String = "mau" 

    val CassandraHostProperty = "CASSANDRA_HOST" 
    val CassandraDefaultHost = "127.0.0.1" 
    val CassandraHost = InetAddress.getByName(sys.env.getOrElse(CassandraHostProperty, CassandraDefaultHost)) 

    val conf = new SparkConf().setAppName(getClass.getSimpleName) 
        .set("spark.cassandra.connection.host", CassandraHost.getHostAddress) 

    lazy val sc = new SparkContext(conf) 
    import com.datastax.spark.connector._ 

    def now = new DateTime(DateTimeZone.UTC) 
    val metricMonth = now.getYear + "-" + now.getMonthOfYear 

    private val mauMonthSB: StringBuilder = new StringBuilder 
    mauMonthSB.append(now.getYear).append("-") 
    if (now.getMonthOfYear < 10) mauMonthSB.append("0") 
    mauMonthSB.append(now.getMonthOfYear).append("-") 
    if (now.getDayOfMonth < 10) mauMonthSB.append("0") 
    mauMonthSB.append(now.getDayOfMonth) 

    private val mauMonth: String = mauMonthSB.toString() 

    val dates = ListBuffer[String]() 
    for (day <- 1 to now.dayOfMonth().getMaximumValue) { 
        val metricDate: StringBuilder = new StringBuilder 
        metricDate.append(now.getYear).append("-") 
        if (now.getMonthOfYear < 10) metricDate.append("0") 
        metricDate.append(now.getMonthOfYear).append("-") 
        if (day < 10) metricDate.append("0") 
        metricDate.append(day) 
        dates += metricDate.toString() 
    } 

    private val metricName: List[String] = List("ms", "md") 
    val appMauAggregate = scala.collection.mutable.Map[String, scala.collection.mutable.Map[UUID, AtomicLong]]() 

    case class MAURecord(month: String, appId: UUID, userId: UUID) extends Serializable 
    case class DUMARecord(metricDate: Date, userId: UUID, metricName: String) extends Serializable 
    case class MAUAggregate(month: String, appId: UUID, total: Long) extends Serializable 

    private val mau = sc.cassandraTable[MAURecord]("analytics", "monthly_active_users") 
        .where("month = ?", metricMonth) 
        .collect() 

    mau.foreach { monthlyActiveUser => 
        val duma = sc.cassandraTable[DUMARecord]("analytics", "daily_user_metric_aggregates") 
            .where("metric_date in ? and user_id = ? and metric_name in ?", dates, monthlyActiveUser.userId, metricName) 
            //.map(_.userId).distinct().collect() 
            .collect() 

        if (duma.length > 0) { // if user has `ms` for the given month 
            if (!appMauAggregate.isDefinedAt(mauMonth)) { 
                appMauAggregate += (mauMonth -> scala.collection.mutable.Map[UUID, AtomicLong]()) 
            } 
            val monthMap: scala.collection.mutable.Map[UUID, AtomicLong] = appMauAggregate(mauMonth) 
            if (!monthMap.isDefinedAt(monthlyActiveUser.appId)) { 
                monthMap += (monthlyActiveUser.appId -> new AtomicLong(0)) 
            } 
            monthMap(monthlyActiveUser.appId).incrementAndGet() 
        } else { 
            println(s"No message_sent in daily_user_metric_aggregates for user: $monthlyActiveUser") 
        } 

    } 
    for ((metricMonth: String, appMauCounts: scala.collection.mutable.Map[UUID, AtomicLong]) <- appMauAggregate) { 
        for ((appId: UUID, total: AtomicLong) <- appMauCounts) { 
            println(s"month: $metricMonth, app_id: $appId, total: $total"); 
            val collection = sc.parallelize(Seq(MAUAggregate(metricMonth.substring(0, 7), appId, total.get()))) 
            collection.saveToCassandra(KeySpace, MauTable, SomeColumns("month", "app_id", "total")) 
        } 
    } 
    sc.stop() 
}

谢谢。

【问题讨论】:

    标签: cassandra apache-spark datastax-enterprise datastax


    【解决方案1】:

    您的解决方案是效率最低的。您通过逐个查找每个键来执行连接,避免任何可能的并行化。

    我从未使用过 Cassandra 连接器,但我知道它会返回 RDD。所以你可以这样做:

    val mau: RDD[(UUID, MAURecord)] = sc
        .cassandraTable[MAURecord]("analytics", "monthly_active_users") 
        .where("month = ?", metricMonth)
        .map(u => u.userId -> u)  // Key by user ID.
    val duma: RDD[(UUID, DUMARecord)] = sc
        .cassandraTable[DUMARecord]("analytics", "daily_user_metric_aggregates") 
        .where("metric_date in ? metric_name in ?", dates, metricName)
        .map(a => a.userId -> a)  // Key by user ID.
    // Count "duma" by key.
    val dumaCounts: RDD[(UUID, Long)] = duma.countByKey
    // Join to "mau". This drops "mau" entries that have no count
    // and "duma" entries that are not present in "mau".
    val joined: RDD[(UUID, (MAURecord, Long))] = mau.join(dumaCounts)
    // Get per-application counts.
    val appCounts: RDD[(UUID, Long)] = joined
        .map { case (u, (mau, count)) => mau.appId -> 1 }
        .countByKey
    

    【讨论】:

      【解决方案2】:
      1. 有一个参数 spark.cassandra.connection.keep_alive_ms 控制连接保持打开的时间。查看文档page

      2. 如果您将 Spark Worker 与 Cassandra 节点放在一起,连接器将利用这一点并适当地创建分区,以便执行程序始终从本地节点获取数据。

      我可以看到您可以在 DUMA 表中进行一些设计改进:metric_date 似乎不是分区键的最佳选择 - 考虑将 (user_id, metric_name) 设置为分区键,因为在这种情况下您不必为查询 - 您只需将 user_id 和 metrics_name 放入 where 子句。此外,您可以向主键添加月份标识符 - 然后,每个分区将仅包含与您希望通过每个查询获取的内容相关的信息。

      无论如何,Spark-Cassandra-Connector 中的 join 功能目前正在实现中(参见this ticket)。

      【讨论】:

        猜你喜欢
        • 2017-08-13
        • 2021-11-21
        • 2015-10-29
        • 1970-01-01
        • 2018-11-28
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-10-22
        相关资源
        最近更新 更多