Spark中获取SaveAsTable与Insert Into操作的记录数方法问询
获取Spark操作Hive各阶段的读写记录数
嘿,我来帮你搞定这两个阶段的记录数统计问题!结合你的Scala代码,我整理了每个环节的具体实现方法:
1. JSON转DataFrame的记录数(你已掌握的方法)
直接用df.count()就能拿到这个阶段的记录数,代码可以这么写:
// read json to DF val df = hiveContext.read.schema(schema1).json(file) val jsonToDfCount = df.count() println(s"JSON转DataFrame的记录数:$jsonToDfCount")
2. DataFrame写入临时表(stg)的记录数
因为saveAsTable是把整个DataFrame完整写入表中(append模式),所以写入前的df.count()就是本次写入stg的记录数。如果担心写入过程中存在数据过滤(比如不符合表结构的行被丢弃),可以通过"写入前后计数差"来验证:
// 先获取写入前stg表的记录数(如果stg已存在) val stgPreCount = hiveContext.sql("select count(*) from stg").head().getLong(0) // DF to Staging df.write.mode("append").saveAsTable("stg") // 计算本次写入的增量记录数 val stgPostCount = hiveContext.sql("select count(*) from stg").head().getLong(0) val writeToStgCount = stgPostCount - stgPreCount println(s"写入临时表stg的记录数:$writeToStgCount")
如果stg是首次创建,stgPreCount默认为0,直接用stgPostCount即可。
3. 从stg插入final表的记录数
这里有两种实用方法,按需选择:
方法一:预查询待插入的记录数(高效推荐)
先执行你的select语句统计待插入数量,再执行insert操作,避免重复计算:
// 先计算符合插入条件的记录数(和insert逻辑完全一致) val insertCount = hiveContext.sql("select count(distinct id) from stg left outer join final on stg.id = final.id where stg.id is not null and final.id is null").head().getLong(0) println(s"待插入final表的记录数:$insertCount") // 执行插入操作 hiveContext.sql("insert into final select distinct columns from stg left outer join final on stg.id = final.id where stg.id is not null and final.id is null")
这里用count(distinct id)是为了和你select语句中的distinct columns逻辑匹配,保证统计数量和实际插入一致。
方法二:通过Spark查询结果获取记录数
利用Spark SQL的执行结果转化为RDD统计,不过这个方法会触发两次计算(统计+插入),数据量大时不推荐:
// 执行插入并获取记录数 val insertQuery = hiveContext.sql("insert into final select distinct columns from stg left outer join final on stg.id = final.id where stg.id is not null and final.id is null") val insertCount = insertQuery.queryExecution.toRdd.count() println(s"插入final表的记录数:$insertCount")
内容的提问来源于stack exchange,提问作者Neha
相关产品推荐
相关产品推荐

