Last active
October 13, 2020 02:15
-
-
Save nsivabalan/a994299cdb502918f0be2234cd573e5f to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| # Licensed to the Apache Software Foundation (ASF) under one | |
| # or more contributor license agreements. See the NOTICE file | |
| # distributed with this work for additional information | |
| # regarding copyright ownership. The ASF licenses this file | |
| # to you under the Apache License, Version 2.0 (the | |
| # "License"); you may not use this file except in compliance | |
| # with the License. You may obtain a copy of the License at | |
| # | |
| # http://www.apache.org/licenses/LICENSE-2.0 | |
| # | |
| # Unless required by applicable law or agreed to in writing, software | |
| # distributed under the License is distributed on an "AS IS" BASIS, | |
| # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | |
| # See the License for the specific language governing permissions and | |
| # limitations under the License. | |
| first_insert: | |
| config: | |
| record_size: 70000 | |
| num_insert_partitions: 1 | |
| repeat_count: 5 | |
| num_records_insert: 1000 | |
| type: InsertNode | |
| deps: none | |
| second_insert: | |
| config: | |
| record_size: 70000 | |
| num_insert_partitions: 1 | |
| repeat_count: 5 | |
| num_records_insert: 10000 | |
| deps: first_insert | |
| type: InsertNode | |
| third_insert: | |
| config: | |
| record_size: 70000 | |
| num_insert_partitions: 1 | |
| repeat_count: 2 | |
| num_records_insert: 300 | |
| deps: second_insert | |
| type: InsertNode | |
| first_upsert: | |
| config: | |
| record_size: 70000 | |
| num_insert_partitions: 1 | |
| num_records_insert: 300 | |
| repeat_count: 5 | |
| num_records_upsert: 100 | |
| num_upsert_partitions: 10 | |
| type: UpsertNode | |
| deps: third_insert | |
| first_hive_sync: | |
| config: | |
| queue_name: "adhoc" | |
| engine: "mr" | |
| type: HiveSyncNode | |
| deps: first_upsert | |
| first_hive_query: | |
| config: | |
| hive_props: | |
| prop2: "set spark.yarn.queue=" | |
| prop3: "set hive.strict.checks.large.query=false" | |
| prop4: "set hive.stats.autogather=false" | |
| hive_queries: | |
| query1: "select count(*) from testdb.table1 group by `_row_key` having count(*) > 1" | |
| result1: 0 | |
| query2: "select count(*) from testdb.table1" | |
| result2: 11600 | |
| type: HiveQueryNode | |
| deps: first_hive_sync | |
| first_validate: | |
| config: | |
| abc: def | |
| type: ValidateDatasetNode | |
| deps: first_hive_query |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 20/10/13 00:32:21 WARN AvroKeyInputFormat: Reader schema was not set. Use AvroJob.setInputKeySchema() if desired. | |
| 20/10/13 00:32:21 WARN AvroKeyInputFormat: Reader schema was not set. Use AvroJob.setInputKeySchema() if desired. | |
| 20/10/13 00:32:22 WARN AvroKeyInputFormat: Reader schema was not set. Use AvroJob.setInputKeySchema() if desired. | |
| 20/10/13 00:32:39 WARN SparkSession$Builder: Using an existing SparkSession; some configuration may not take effect. | |
| 20/10/13 00:32:39 WARN ValidateDatasetNode: ValidateDataset Node: Input path /user/hive/warehouse/hudi-integ-test-suite/output/../input/*/*, hudi path /user/hive/warehouse/hudi-integ-test-suite/output/*/*/* | |
| 20/10/13 00:32:42 WARN DefaultSource: Loading Base File Only View. | |
| 20/10/13 00:32:48 ERROR HoodieTestSuiteJob: Failed to run Test Suite | |
| java.util.concurrent.ExecutionException: java.lang.NoSuchFieldError: HIVE_STATS_JDBC_TIMEOUT | |
| at java.util.concurrent.FutureTask.report(FutureTask.java:122) | |
| at java.util.concurrent.FutureTask.get(FutureTask.java:206) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.execute(DagScheduler.java:109) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.schedule(DagScheduler.java:68) | |
| at org.apache.hudi.integ.testsuite.HoodieTestSuiteJob.runTestSuite(HoodieTestSuiteJob.java:137) | |
| at org.apache.hudi.integ.testsuite.HoodieTestSuiteJob.main(HoodieTestSuiteJob.java:117) | |
| at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) | |
| at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) | |
| at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) | |
| at java.lang.reflect.Method.invoke(Method.java:498) | |
| at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) | |
| at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:845) | |
| at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:161) | |
| at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:184) | |
| at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:86) | |
| at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:920) | |
| at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:929) | |
| at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala) | |
| Caused by: java.lang.NoSuchFieldError: HIVE_STATS_JDBC_TIMEOUT | |
| at org.apache.spark.sql.hive.HiveUtils$.formatTimeVarsForHiveClient(HiveUtils.scala:204) | |
| at org.apache.spark.sql.hive.HiveUtils$.newClientForMetadata(HiveUtils.scala:285) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.client$lzycompute(HiveExternalCatalog.scala:66) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.client(HiveExternalCatalog.scala:65) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply$mcZ$sp(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.withClient(HiveExternalCatalog.scala:97) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.databaseExists(HiveExternalCatalog.scala:214) | |
| at org.apache.spark.sql.internal.SharedState.externalCatalog$lzycompute(SharedState.scala:114) | |
| at org.apache.spark.sql.internal.SharedState.externalCatalog(SharedState.scala:102) | |
| at org.apache.spark.sql.internal.SharedState.globalTempViewManager$lzycompute(SharedState.scala:141) | |
| at org.apache.spark.sql.internal.SharedState.globalTempViewManager(SharedState.scala:136) | |
| at org.apache.spark.sql.hive.HiveSessionStateBuilder$$anonfun$2.apply(HiveSessionStateBuilder.scala:55) | |
| at org.apache.spark.sql.hive.HiveSessionStateBuilder$$anonfun$2.apply(HiveSessionStateBuilder.scala:55) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.globalTempViewManager$lzycompute(SessionCatalog.scala:91) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.globalTempViewManager(SessionCatalog.scala:91) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.isTemporaryTable(SessionCatalog.scala:736) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.isRunningDirectlyOnFiles(Analyzer.scala:747) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.resolveRelation(Analyzer.scala:681) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$$anonfun$apply$8.applyOrElse(Analyzer.scala:713) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$$anonfun$apply$8.applyOrElse(Analyzer.scala:706) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$apply$1.apply(AnalysisHelper.scala:90) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$apply$1.apply(AnalysisHelper.scala:90) | |
| at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:70) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:89) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$class.resolveOperatorsUp(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolveOperatorsUp(LogicalPlan.scala:29) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$4.apply(TreeNode.scala:329) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:187) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:327) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$class.resolveOperatorsUp(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolveOperatorsUp(LogicalPlan.scala:29) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.apply(Analyzer.scala:706) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.apply(Analyzer.scala:652) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:87) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:84) | |
| at scala.collection.LinearSeqOptimized$class.foldLeft(LinearSeqOptimized.scala:124) | |
| at scala.collection.immutable.List.foldLeft(List.scala:84) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:84) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:76) | |
| at scala.collection.immutable.List.foreach(List.scala:392) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor.execute(RuleExecutor.scala:76) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.org$apache$spark$sql$catalyst$analysis$Analyzer$$executeSameContext(Analyzer.scala:127) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.execute(Analyzer.scala:121) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$$anonfun$executeAndCheck$1.apply(Analyzer.scala:106) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$$anonfun$executeAndCheck$1.apply(Analyzer.scala:105) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.markInAnalyzer(AnalysisHelper.scala:201) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.executeAndCheck(Analyzer.scala:105) | |
| at org.apache.spark.sql.execution.QueryExecution.analyzed$lzycompute(QueryExecution.scala:57) | |
| at org.apache.spark.sql.execution.QueryExecution.analyzed(QueryExecution.scala:55) | |
| at org.apache.spark.sql.execution.QueryExecution.assertAnalyzed(QueryExecution.scala:47) | |
| at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:78) | |
| at org.apache.spark.sql.SparkSession.sql(SparkSession.scala:642) | |
| at org.apache.hudi.integ.testsuite.dag.nodes.ValidateDatasetNode.execute(ValidateDatasetNode.java:64) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.executeNode(DagScheduler.java:131) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.lambda$execute$0(DagScheduler.java:101) | |
| at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) | |
| at java.util.concurrent.FutureTask.run(FutureTask.java:266) | |
| at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) | |
| at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) | |
| at java.lang.Thread.run(Thread.java:748) | |
| Exception in thread "main" org.apache.hudi.exception.HoodieException: Failed to run Test Suite | |
| at org.apache.hudi.integ.testsuite.HoodieTestSuiteJob.runTestSuite(HoodieTestSuiteJob.java:141) | |
| at org.apache.hudi.integ.testsuite.HoodieTestSuiteJob.main(HoodieTestSuiteJob.java:117) | |
| at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) | |
| at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) | |
| at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) | |
| at java.lang.reflect.Method.invoke(Method.java:498) | |
| at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) | |
| at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:845) | |
| at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:161) | |
| at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:184) | |
| at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:86) | |
| at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:920) | |
| at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:929) | |
| at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala) | |
| Caused by: java.util.concurrent.ExecutionException: java.lang.NoSuchFieldError: HIVE_STATS_JDBC_TIMEOUT | |
| at java.util.concurrent.FutureTask.report(FutureTask.java:122) | |
| at java.util.concurrent.FutureTask.get(FutureTask.java:206) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.execute(DagScheduler.java:109) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.schedule(DagScheduler.java:68) | |
| at org.apache.hudi.integ.testsuite.HoodieTestSuiteJob.runTestSuite(HoodieTestSuiteJob.java:137) | |
| ... 13 more | |
| Caused by: java.lang.NoSuchFieldError: HIVE_STATS_JDBC_TIMEOUT | |
| at org.apache.spark.sql.hive.HiveUtils$.formatTimeVarsForHiveClient(HiveUtils.scala:204) | |
| at org.apache.spark.sql.hive.HiveUtils$.newClientForMetadata(HiveUtils.scala:285) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.client$lzycompute(HiveExternalCatalog.scala:66) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.client(HiveExternalCatalog.scala:65) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply$mcZ$sp(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog$$anonfun$databaseExists$1.apply(HiveExternalCatalog.scala:215) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.withClient(HiveExternalCatalog.scala:97) | |
| at org.apache.spark.sql.hive.HiveExternalCatalog.databaseExists(HiveExternalCatalog.scala:214) | |
| at org.apache.spark.sql.internal.SharedState.externalCatalog$lzycompute(SharedState.scala:114) | |
| at org.apache.spark.sql.internal.SharedState.externalCatalog(SharedState.scala:102) | |
| at org.apache.spark.sql.internal.SharedState.globalTempViewManager$lzycompute(SharedState.scala:141) | |
| at org.apache.spark.sql.internal.SharedState.globalTempViewManager(SharedState.scala:136) | |
| at org.apache.spark.sql.hive.HiveSessionStateBuilder$$anonfun$2.apply(HiveSessionStateBuilder.scala:55) | |
| at org.apache.spark.sql.hive.HiveSessionStateBuilder$$anonfun$2.apply(HiveSessionStateBuilder.scala:55) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.globalTempViewManager$lzycompute(SessionCatalog.scala:91) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.globalTempViewManager(SessionCatalog.scala:91) | |
| at org.apache.spark.sql.catalyst.catalog.SessionCatalog.isTemporaryTable(SessionCatalog.scala:736) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.isRunningDirectlyOnFiles(Analyzer.scala:747) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.resolveRelation(Analyzer.scala:681) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$$anonfun$apply$8.applyOrElse(Analyzer.scala:713) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$$anonfun$apply$8.applyOrElse(Analyzer.scala:706) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$apply$1.apply(AnalysisHelper.scala:90) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$apply$1.apply(AnalysisHelper.scala:90) | |
| at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:70) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:89) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$class.resolveOperatorsUp(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolveOperatorsUp(LogicalPlan.scala:29) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1$$anonfun$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode$$anonfun$4.apply(TreeNode.scala:329) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:187) | |
| at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:327) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:87) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$$anonfun$resolveOperatorsUp$1.apply(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$class.resolveOperatorsUp(AnalysisHelper.scala:86) | |
| at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolveOperatorsUp(LogicalPlan.scala:29) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.apply(Analyzer.scala:706) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$ResolveRelations$.apply(Analyzer.scala:652) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:87) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1$$anonfun$apply$1.apply(RuleExecutor.scala:84) | |
| at scala.collection.LinearSeqOptimized$class.foldLeft(LinearSeqOptimized.scala:124) | |
| at scala.collection.immutable.List.foldLeft(List.scala:84) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:84) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor$$anonfun$execute$1.apply(RuleExecutor.scala:76) | |
| at scala.collection.immutable.List.foreach(List.scala:392) | |
| at org.apache.spark.sql.catalyst.rules.RuleExecutor.execute(RuleExecutor.scala:76) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.org$apache$spark$sql$catalyst$analysis$Analyzer$$executeSameContext(Analyzer.scala:127) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.execute(Analyzer.scala:121) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$$anonfun$executeAndCheck$1.apply(Analyzer.scala:106) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer$$anonfun$executeAndCheck$1.apply(Analyzer.scala:105) | |
| at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.markInAnalyzer(AnalysisHelper.scala:201) | |
| at org.apache.spark.sql.catalyst.analysis.Analyzer.executeAndCheck(Analyzer.scala:105) | |
| at org.apache.spark.sql.execution.QueryExecution.analyzed$lzycompute(QueryExecution.scala:57) | |
| at org.apache.spark.sql.execution.QueryExecution.analyzed(QueryExecution.scala:55) | |
| at org.apache.spark.sql.execution.QueryExecution.assertAnalyzed(QueryExecution.scala:47) | |
| at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:78) | |
| at org.apache.spark.sql.SparkSession.sql(SparkSession.scala:642) | |
| at org.apache.hudi.integ.testsuite.dag.nodes.ValidateDatasetNode.execute(ValidateDatasetNode.java:64) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.executeNode(DagScheduler.java:131) | |
| at org.apache.hudi.integ.testsuite.dag.scheduler.DagScheduler.lambda$execute$0(DagScheduler.java:101) | |
| at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) | |
| at java.util.concurrent.FutureTask.run(FutureTask.java:266) | |
| at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) | |
| at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) | |
| at java.lang.Thread.run(Thread.java:748) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| public class ValidateDatasetNode extends DagNode<Boolean> { | |
| private static Logger log = LoggerFactory.getLogger(ValidateDatasetNode.class); | |
| public ValidateDatasetNode(Config config) { | |
| this.config = config; | |
| } | |
| @Override | |
| public void execute(ExecutionContext context) throws Exception { | |
| SparkConf sparkConf = new SparkConf().setAppName("ValidateApp").setMaster("local"); | |
| SparkSession spark = SparkSession | |
| .builder() | |
| .config(sparkConf) | |
| .getOrCreate(); | |
| String inputPath = context.getHoodieTestSuiteWriter().getCfg().targetBasePath + "/../input/*/*"; | |
| String hudiPath = context.getHoodieTestSuiteWriter().getCfg().targetBasePath + "/*/*/*"; | |
| log.warn("ValidateDataset Node: Input path " + inputPath + ", hudi path " + hudiPath); | |
| Dataset<Row> inputDf = spark.read().format("avro").load(inputPath); | |
| Dataset<Row> hudiDf = spark.read().format("hudi").load(hudiPath); | |
| Dataset<Row> trimmedDf = hudiDf.drop(HoodieRecord.COMMIT_TIME_METADATA_FIELD).drop(HoodieRecord.COMMIT_SEQNO_METADATA_FIELD).drop(HoodieRecord.RECORD_KEY_METADATA_FIELD) | |
| .drop(HoodieRecord.PARTITION_PATH_METADATA_FIELD).drop(HoodieRecord.FILENAME_METADATA_FIELD); | |
| if (inputDf.except(trimmedDf).count() != 0) { | |
| log.error("Data set validation failed. Total count in hudi " + trimmedDf.count() + ", input df count " + inputDf.count()); | |
| throw new AssertionError("Hudi contents does not match contents input data. "); | |
| } | |
| Dataset<Row> cowDf = spark.sql("SELECT * FROM testdb.table1"); | |
| Dataset<Row> trimmedCowDf = cowDf.drop(HoodieRecord.COMMIT_TIME_METADATA_FIELD).drop(HoodieRecord.COMMIT_SEQNO_METADATA_FIELD).drop(HoodieRecord.RECORD_KEY_METADATA_FIELD) | |
| .drop(HoodieRecord.PARTITION_PATH_METADATA_FIELD).drop(HoodieRecord.FILENAME_METADATA_FIELD); | |
| if (inputDf.except(trimmedCowDf).count() != 0) { | |
| log.error("Data set validation failed for COW table. Total count in hudi " + trimmedCowDf.count() + ", input df count " + inputDf.count()); | |
| throw new AssertionError("Hudi contents does not match contents input data. "); | |
| } | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment