New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
MutationScan seems to be doing the wrong thing #755
Conversation
@@ -515,8 +515,8 @@ object GroupBy { | |||
// Generate mutation Df if required, align the columns with inputDf so no additional schema is needed by aggregator. | |||
val mutationSources = groupByConf.sources.toScala.filter { _.isSetEntities } | |||
val mutationsColumnOrder = inputDf.columns ++ Constants.MutationFields.map(_.name) | |||
val mutationDf = | |||
if (mutationScan && groupByConf.inferredAccuracy == api.Accuracy.TEMPORAL && mutationSources.nonEmpty) { | |||
def mutationDfFn(): DataFrame = { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This makes mutationDf eval lazy.
So that it is only touched when temporal entities code is run
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nice!
c0f44f6
to
6f56099
Compare
@@ -325,4 +329,12 @@ object Extensions { | |||
} | |||
} | |||
} | |||
|
|||
implicit class TupleToJMapOps[K, V](tuples: Iterator[(K, V)]) { | |||
def toJMap: util.Map[K, V] = { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Will it be more performant than the Scala Map?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Scala map somehow is not serializable :-/
// events | entities | snapshot => right part tables are not aligned - so scan by leftTimeRange | ||
// events | entities | temporal => right part tables are aligned - so scan by leftRange | ||
// entities | entities | snapshot => right part tables are aligned - so scan by leftRange | ||
val rightRange = if (joinConf.left.dataModel == Events && joinPart.groupBy.inferredAccuracy == Accuracy.SNAPSHOT) { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
thanks for the comments above
@@ -90,7 +91,8 @@ object SparkSessionBuilder { | |||
} | |||
val spark = builder.getOrCreate() | |||
// disable log spam | |||
spark.sparkContext.setLogLevel("ERROR") | |||
// spark.sparkContext.setLogLevel("ERROR") |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
do you want to keep this line?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think there is an issue with logs right now - will follow up PR to fix the logs issue
@@ -117,7 +119,7 @@ object SparkSessionBuilder { | |||
} | |||
val spark = builder.getOrCreate() | |||
// disable log spam | |||
spark.sparkContext.setLogLevel("ERROR") | |||
// spark.sparkContext.setLogLevel("ERROR") |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
do you want to keep this line?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
There is something wrong with our logs - they just disappear in the tests
|
||
val DefaultWarehouseDir = new File("/tmp/chronon/spark-warehouse") | ||
private val warehouseId = java.util.UUID.randomUUID().toString.takeRight(6) |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
would it solve the potential racing condition in the CI test?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think it does!
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Thanks for cleaning up and removing unused mutationScan
Summary
when joinSource left is events & right is entities, mutationScan is being set to false, but we still enter temporalEntities case as we should. This PR disables mutation scan all together.
Test Plan
Checklist
Reviewers
@ezvz @better365