Tuesday, April 17, 2018
Solution For "Error: java.io.IOException: org.apache.parquet.io.ParquetDecodingException: Can not read value at 0 in block -1 in file" in HDFS
This is due to the Schema Revolution feature of parquet files and the column name of that parquet file may have changed before. so try to apply `set parquet.column.index.access=true` and the issue is solved.
Monday, April 16, 2018
Hive的mr作业产生很多小文件或空文件的解决方案
set hive.merge.mapfiles=true;
set hive.merge.mapredfiles=true;
set hive.merge.smallfiles.avgsize=16000000;
set hive.merge.size.per.task=67108864;
引用一段具体的scenario description:
By default hive.merge.smallfiles.avgsize=16000000 and hive.merge.size.per.task=256000000, so if the average file size is about 17MB, the merge job will not be triggered. Sometimes if we really want only 1 file being generated in the end, we need to increase hive.merge.smallfiles.avgsize to large enough to trigger the merge; and also you need to increase hive.merge.size.per.task to the get the needed number of files in the end.
REFERENCE:
HiveOnSpark系列:spark jobs partition数量调优问题
在
HiveOnSpark上跑Spark application时,发现部分stage对应的task数量很少,导致full-gc严重。例如下图中的stage只有17个tasks在跑。即使设置了spark.sql.shuffle.partitions=1201和spark.default.parallelism=1202也没有用,依旧是17个。
通过trace spark源码发现
mapred.reduce.tasks参数虽然已经deprecated,但会优先spark.sql.shuffle.partitions设置到环境变量中。(spark-2.0:SetCommand:44)
同时,根据同事给的hadoop参数优化项里
hive.exec.reducers.bytes.per.reducer,把两者按如下配置后,4个stage启动的tasks数量分别从82/33/17/526变为了82/1201/1201/1201,大大增加了partition数量并as expected. 由此可见,HiveOnSpark时,Spark的execution plan应该还是走了Hive的逻辑,所以部分hadoop相关的参数会主导spark的执行计划的生成逻辑和结果
set hive.exec.reducers.bytes.per.reducer=67108864;
set mapred.reduce.tasks=1201;
P.S.
Hive version: 0.23
Spark version: 2.0.3
REFERENCE:
- https://cwiki.apache.org/confluence/display/Hive/Configuration+Properties#ConfigurationProperties-hive.exec.reducers.bytes.per.reducer
Hive Hook遇到的坑儿和解决思路
问题要从Hive Execution PreHook说起。有个需求是根据SQL中用到的table/partition总大小来判断用MR还是Spark。但在执行的时候,虽然确定prehook中已经将
hive.execution.engine设置为spark,但执行的时候还是使用了默认的mr。
在hive源代码里全局搜
hive.exec.pre.hooks, 可以找到入口在HiveConf的enum类中:PREEXECHOOKS("hive.exec.pre.hooks", "",
"Comma-separated list of pre-execution hooks to be invoked for each statement. \n" +
"A pre-execution hook is specified as the name of a Java class which implements the \n" +
"org.apache.hadoop.hive.ql.hooks.ExecuteWithHookContext interface."),
在
Driver类里搜索“PREEXECHOOKS”,可以发现调用的方法在Driver.execute(). 在此处抛出stacktrace(手动throw exception然后catch打印stacktrace即可)应该能获得类似如下的信息(测试时是在Optimizer打印的stacktrace):2018-04-12T19:55:21,535 INFO [1df5090e-6c7f-4751-8979-d4fd3ef5024e HiveServer2-Handler-Pool: Thread-92] optimize
r.Optimizer: [FLAG_15_1] stack trace java.lang.RuntimeException: t
at org.apache.hadoop.hive.ql.optimizer.Optimizer.initialize(Optimizer.java:66)
at org.apache.hadoop.hive.ql.parse.SemanticAnalyzer.analyzeInternal(SemanticAnalyzer.java:11246)
at org.apache.hadoop.hive.ql.parse.CalcitePlanner.analyzeInternal(CalcitePlanner.java:286)
at org.apache.hadoop.hive.ql.parse.BaseSemanticAnalyzer.analyze(BaseSemanticAnalyzer.java:259)
at org.apache.hadoop.hive.ql.Driver.compile(Driver.java:814)
at org.apache.hadoop.hive.ql.Driver.compileInternal(Driver.java:1286)
at org.apache.hadoop.hive.ql.Driver.compileAndRespond(Driver.java:1265)
at org.apache.hive.service.cli.operation.SQLOperation.prepare(SQLOperation.java:204)
at org.apache.hive.service.cli.operation.SQLOperation.runInternal(SQLOperation.java:290)
at org.apache.hive.service.cli.operation.Operation.run(Operation.java:320)
at org.apache.hive.service.cli.session.HiveSessionImpl.executeStatementAsyncntInternal(HiveSessionImpl.java:530)
at org.apache.hive.service.cli.session.HiveSessionImpl.executeStatementAsync(HiveSessionImpl.java:517)
移步到Driver.compile:814, 找到关于semantic analyzer相关逻辑:
BaseSemanticAnalyzer sem = SemanticAnalyzerFactory.get(queryState, tree);
List<HiveSemanticAnalyzerHook> saHooks =
getHooks(HiveConf.ConfVars.SEMANTIC_ANALYZER_HOOK,
HiveSemanticAnalyzerHook.class);
// Flush the metastore cache. This assures that we don't pick up objects from a previous
// query running in this same thread. This has to be done after we get our semantic
// analyzer (this is when the connection to the metastore is made) but before we analyze,
// because at that point we need access to the objects.
Hive.get().getMSC().flushCache();
// Do semantic analysis and plan generation
if (saHooks != null && !saHooks.isEmpty()) {
HiveSemanticAnalyzerHookContext hookCtx = new HiveSemanticAnalyzerHookContextImpl();
hookCtx.setConf(conf);
hookCtx.setUserName(userName);
hookCtx.setIpAddress(SessionState.get().getUserIpAddress());
hookCtx.setCommand(command);
hookCtx.setHiveOperation(queryState.getHiveOperation());
hookCtx.setQueryState(queryState);
hookCtx.setContext(ctx);
for (HiveSemanticAnalyzerHook hook : saHooks) {
tree = hook.preAnalyze(hookCtx, tree);
}
sem.analyze(tree, ctx);
hookCtx.update(sem);
for (HiveSemanticAnalyzerHook hook : saHooks) {
hook.postAnalyze(hookCtx, sem.getAllRootTasks());
}
} else {
sem.analyze(tree, ctx);
}
此时可以判断得出的是,在execution.prehook阶段,本身任务是mr还是spark早已确定,通过set conf肯定是无法实现的。那么此过程需要在semantic_analyzer.prehook阶段完成。但由于在sementic analyze步骤之前没有语法解析,所以没有input table/partition信息。这是解决的思路之一是,ch同构hive源码,将initial
sem.analyze(tree, ctx)所需的两个参数传给semantic_analyzer.prehook,在里面先执行analyze()方法拿到input信息,从hive metastore拿到对应大小后再set conf从而决定使用mr还是sparkHive: The following columns have types incompatible with the existing columns
Since Hive can't feel compatibility of complicated data structure like struct, so if you intend to change a column of type array<struct<a:int, b:int>> to array<struct<a:int, b:int, c:int>>, it will complains error: "The following columns have types incompatible with the existing columns".
It's easy to depress the above check by setting `set hive.metastore.disallow.incompatible.col.type.changes=false;` and khalas.
It's easy to depress the above check by setting `set hive.metastore.disallow.incompatible.col.type.changes=false;` and khalas.
Sunday, April 15, 2018
Spark Parameter Tuning Cases Due to FullGC and Data Skew
The default Spark setting is as follows:
> spark.executor.memory 8g
> spark.executor.cores 4
> spark.dynamicAllocation.enabled true
> spark.dynamicAllocation.initialExecutors 1
> spark.dynamicAllocation.maxExecutors 200
> spark.dynamicAllocation.minExecutors 1
when running SQL as below, it shows that GC time occupies virtually as much as the running time, which, apparently, is due to a general memory issue. Overall, it took 26min to finish.
SELECT a, b, c, d, e,
avg(f/g) AS x1,
percentile(cast(f/g AS bigint), array(0.5, 0.95)) AS x2
FROM (
SELECT a, b, c, d, e,
h, sum(i) AS f, sum(1) AS g
FROM tbl_name
WHERE p_date = '20180303'
GROUP BY a, b, c, d, e,
h
) AS TZ
GROUP BY a, b, c, d, e WITH CUBE
SORT BY a, b DESC, c, d, e;
after updating the following parameters, execution time reduces to 16min and GC time becomes much less dominant.
> set spark.executor.extraJavaOptions=-XX:+UseG1GC;
> set spark.yarn.executor.memoryOverhead=8g;
> set spark.executor.memory=24g;
Yet the following issue is data skew. there's stragglers running for a long time whereas the others have already finished quickly.
This is because of the default value for both `spark.sql.shuffle.partitions` and `spark.default.parallelism`, a small setting value and large amount of data is more likely to lead to data skew and stragglers issue. After updating as follows, execution time goes down to 705s in total.
> set spark.sql.shuffle.partitions=2467;
> set spark.default.parallelism=127;
REFERENCE:
1. https://stackoverflow.com/questions/45704156/what-is-the-difference-between-spark-sql-shuffle-partitions-and-spark-default-pa
2. https://www.jianshu.com/p/06b67a3c61a9
3. https://databricks.com/blog/2015/05/28/tuning-java-garbage-collection-for-spark-applications.html
4. https://www.ibm.com/support/knowledgecenter/en/SS3H8V_1.1.0/com.ibm.izoda.v1r1.azka100/topics/azkic_t_configmemcpu.htm
5. http://spark.apache.org/docs/latest/tuning.html#garbage-collection-tuning
6. http://www.iteye.com/news/31303
> spark.executor.memory 8g
> spark.executor.cores 4
> spark.dynamicAllocation.enabled true
> spark.dynamicAllocation.initialExecutors 1
> spark.dynamicAllocation.maxExecutors 200
> spark.dynamicAllocation.minExecutors 1
when running SQL as below, it shows that GC time occupies virtually as much as the running time, which, apparently, is due to a general memory issue. Overall, it took 26min to finish.
SELECT a, b, c, d, e,
avg(f/g) AS x1,
percentile(cast(f/g AS bigint), array(0.5, 0.95)) AS x2
FROM (
SELECT a, b, c, d, e,
h, sum(i) AS f, sum(1) AS g
FROM tbl_name
WHERE p_date = '20180303'
GROUP BY a, b, c, d, e,
h
) AS TZ
GROUP BY a, b, c, d, e WITH CUBE
SORT BY a, b DESC, c, d, e;
after updating the following parameters, execution time reduces to 16min and GC time becomes much less dominant.
> set spark.executor.extraJavaOptions=-XX:+UseG1GC;
> set spark.yarn.executor.memoryOverhead=8g;
> set spark.executor.memory=24g;
Yet the following issue is data skew. there's stragglers running for a long time whereas the others have already finished quickly.
This is because of the default value for both `spark.sql.shuffle.partitions` and `spark.default.parallelism`, a small setting value and large amount of data is more likely to lead to data skew and stragglers issue. After updating as follows, execution time goes down to 705s in total.
> set spark.sql.shuffle.partitions=2467;
> set spark.default.parallelism=127;
REFERENCE:
1. https://stackoverflow.com/questions/45704156/what-is-the-difference-between-spark-sql-shuffle-partitions-and-spark-default-pa
2. https://www.jianshu.com/p/06b67a3c61a9
3. https://databricks.com/blog/2015/05/28/tuning-java-garbage-collection-for-spark-applications.html
4. https://www.ibm.com/support/knowledgecenter/en/SS3H8V_1.1.0/com.ibm.izoda.v1r1.azka100/topics/azkic_t_configmemcpu.htm
5. http://spark.apache.org/docs/latest/tuning.html#garbage-collection-tuning
6. http://www.iteye.com/news/31303
Monday, April 2, 2018
third-party classpath issue in Intellij
Scenario:
There are two maven projects, A and B respectively, and B has dependency in pom.xml from A. After adding new class in project A and have completed `mvn clean install`, it could still not be found from project B.
Solution:
This is due to a bug at the maven plugin in Intellij, which could be forcibly refreshed via right-clicking on the project in B -> Maven -> Reimport.
There are two maven projects, A and B respectively, and B has dependency in pom.xml from A. After adding new class in project A and have completed `mvn clean install`, it could still not be found from project B.
Solution:
This is due to a bug at the maven plugin in Intellij, which could be forcibly refreshed via right-clicking on the project in B -> Maven -> Reimport.
Subscribe to:
Posts (Atom)





