Spark Sql 之 cacheTable

1. 前言

spark sql中使用DataFrame/DataSet来抽象表示结构化数据(关系数据库中的table),DataSet上支持和RDD类似的操作,和RDD上的操作生成新的RDD一样,DataSet上的操作生成新的DataSet来表示新的数据抽象。最终DataSet上的这些操作经过:
logical plan -> analyzed logical plan -> optimized logical pan -> physical plan -> rdd dag的转化提交rdd 运行。这里plan(执行计划)就是DataSet上的转换操作,一个DataSet也就是对应一个logical plan生成的数据。

cacheTable也就是缓存DataSet抽象表示的数据,也就是DataSet的plan生成的数据。

2. cacheTable

从上面的介绍可以看出DataSet只是数据的抽象,它描述了从数据源头开始经过怎样的执行计划(plan)才能得到当前的DataSet表示的真实数据,也就是必须等到执行计划提交spark job运行结束后才能得到数据。spark实现cacheTable时,并没有立即提交table(DataSet)对应的plan去运行,然后得到运行结果数据去缓存,而是采用一种lazy模式:最终在DataSet上调用一些触发任务提交的方法时(类似RDD的action操作),发现plan对应的抽象语法树中发现子树是表缓存plan,如果这个时候数据已经缓存了,直接使用缓存的数据,没有则触发缓存表的plan去执行,然后采用按列缓存的方式缓存数据。

看看代码实现:
调用SQLContext # cacheTable(tableName : String)最终会走到下面的调用:

// query 即缓存的Dataset
// storageLevel 可以使用memory和disk缓存
def cacheQuery(
      query: Dataset[_],
      tableName: Option[String] = None,
      storageLevel: StorageLevel = MEMORY_AND_DISK): Unit = writeLock {
    // 拿到dataset的plan
    val planToCache = query.logicalPlan
    if (lookupCachedData(planToCache).nonEmpty) {
      logWarning("Asked to cache already cached data.")
    } else {
      val sparkSession = query.sparkSession
      // 缓存是建立一个plan到InMemoryRelation的映射
      cachedData.add(CachedData(
        planToCache,
        InMemoryRelation(
          sparkSession.sessionState.conf.useCompression,
          sparkSession.sessionState.conf.columnBatchSize,
          storageLevel,
          sparkSession.sessionState.executePlan(planToCache).executedPlan,
          tableName)))
    }
  }

上面InMemoryRelation是执行计划中一个节点,当出现select * from table_a语句(或者任何逻辑执行计划中有table_a出现),假设table_a被缓存了,那么这条语句生成的逻辑执行计划中,table_a对应的Relation节点会被执行计划优化器(optimizer)替换成InMemoryRelation。

InMemoryRelation构造参数:

  • columnBatchSize,后面会提到,table缓存是按列缓存的,然后数据又被按行分为一个个batch,这个参数用来控制一个batch里行数,通过配置项spark.sql.inMemoryColumnarStorage.batchSize设置,默认是10000行。

下面是InMemoryRelation中和缓存相关的代码:


  private def buildBuffers(): Unit = {
    // output输出的是Seq[Attribute],也就是表的schema,包含所有列名,列类型等信息
    // child也就是缓存的dataset对应的plan
    val output = child.output
    // 调用逻辑执行计划的execute返回的是RDD[InternalRow],返回RDD是整个执行计划分析的最后一步了,接下来就rdd的提交运行。那么这个rdd也就是dataset表示的数据的rdd形式的抽象。
   // 这里在rdd上调用mapPartitionsInternal,实现的是将遍历每一行数据,然后按列缓存。
   // 这个地方返回新的RDD的数据类型CachedBatch,CachedBatch是一个batch内若干行上的按列缓存。
    val cached = child.execute().mapPartitionsInternal { rowIterator =>
      new Iterator[CachedBatch] {
        def next(): CachedBatch = {
          // 按照每一列的类型生成ColumnBuilder,内部使用数组来保存列数据
          val columnBuilders = output.map { attribute =>
            ColumnBuilder(attribute.dataType, batchSize, attribute.name, useCompression)
          }.toArray

          var rowCount = 0
          var totalSize = 0L
         // 遍历每一行数据,控制当前batch行数 rowCount不超过batchSize,且同时batch中数据大小不超过MAX_BATCH_SIZE_IN_BYTE(4MB)
          while (rowIterator.hasNext && rowCount < batchSize
            && totalSize < ColumnBuilder.MAX_BATCH_SIZE_IN_BYTE) {
            val row = rowIterator.next()

            assert(
              row.numFields == columnBuilders.length,
              s"Row column number mismatch, expected ${output.size} columns, " +
                s"but got ${row.numFields}." +
                s"\nRow content: $row")

            var i = 0
            totalSize = 0
            while (i < row.numFields) {
              columnBuilders(i).appendFrom(row, i)
              totalSize += columnBuilders(i).columnStats.sizeInBytes
              i += 1
            }
            rowCount += 1
          }

          batchStats.add(totalSize)

          val stats = InternalRow.fromSeq(columnBuilders.map(_.columnStats.collectedStatistics)
            .flatMap(_.values))
          CachedBatch(rowCount, columnBuilders.map { builder =>
            JavaUtils.bufferToArray(builder.build())
          }, stats)
        }

        def hasNext: Boolean = rowIterator.hasNext
      }
      // 调用persist缓存RDD,所以cacheTable最终还是调用rdd的缓存接口完成缓存的
    }.persist(storageLevel)

    cached.setName(
      tableName.map(n => s"In-memory table $n")
        .getOrElse(StringUtils.abbreviate(child.toString, 1024)))
  
    _cachedColumnBuffers = cached
  }

上面代码可以看出cacheTable实际上还是通过cache rdd实现的。上面InMemoryRelation只是逻辑执行计划中一个节点,逻辑执行计划需要转换成物理执行计划,再转换成RDD dag才能执行,加上spark中RDD的计算是lazy模式的,所以上面的缓存rdd并没有提交运行,所以数据还没有缓存下来。

真正缓存还得看InMemoryRelation所在的执行计划真正提交后,这个缓存rdd被计算,数据才会被缓存在内存中。

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 212,222评论 6 493
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 90,455评论 3 385
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 157,720评论 0 348
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 56,568评论 1 284
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 65,696评论 6 386
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 49,879评论 1 290
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,028评论 3 409
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 37,773评论 0 268
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,220评论 1 303
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 36,550评论 2 327
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 38,697评论 1 341
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,360评论 4 332
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,002评论 3 315
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 30,782评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,010评论 1 266
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 46,433评论 2 360
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 43,587评论 2 350

推荐阅读更多精彩内容