@@ -23,8 +23,8 @@ case class WorkFlowUsageMetricsAlgoOutput(event_date: Date, total_content_play_s
2323
2424object UpdateWorkFlowUsageMetricsModel extends IBatchModelTemplate [DerivedEvent , DerivedEvent , WorkFlowUsageMetricsAlgoOutput , WorkFlowUsageMetricsAlgoOutput ] with Serializable {
2525
26- val className = " org.ekstep.analytics.updater.UpdateMetrics "
27- override def name : String = " UpdateMetrics "
26+ val className = " org.ekstep.analytics.updater.UpdateWorkFlowUsageMetricsModel "
27+ override def name : String = " UpdateWorkFlowUsageMetricsModel "
2828
2929 /**
3030 * preProcess which will fetch the `ME_WORKFLOW_USAGE_SUMMARY` Event data from the Cassandra Database.
@@ -62,39 +62,40 @@ object UpdateWorkFlowUsageMetricsModel extends IBatchModelTemplate[DerivedEvent,
6262 val SESSION = " session"
6363 val ALL = " all"
6464 }
65-
66- val eventDate = data.first().syncts
67-
68- val totalContentPlaySessions = data.filter(x =>
69- x.dimensions.mode.getOrElse(" " ).equalsIgnoreCase(_constant.PLAY ) &&
70- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.CONTENT )).map { f =>
71- val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
72- eksMap.getOrElse(" total_sessions" , 0 ).asInstanceOf [Number ].longValue()
73- }.sum().toLong
74-
75- val totalTimeSpent = data.filter(x =>
76- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
77- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
78- val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
79- eksMap.getOrElse(" total_ts" , 0.0 ).asInstanceOf [Number ].doubleValue()
80- }.sum()
81-
82- val totalInteractions = data.filter(x =>
83- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
84- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
85- val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
86- eksMap.getOrElse(" total_interactions" , 0 ).asInstanceOf [Number ].longValue()
87- }.sum().toLong
88-
89- val totalPageviews = data.filter(x =>
90- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
91- x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
92- val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
93- eksMap.getOrElse(" total_pageviews_count" , 0 ).asInstanceOf [Number ].longValue()
94- }.sum().toLong
95-
96- sc.parallelize(Array (WorkFlowUsageMetricsAlgoOutput (new Date (eventDate), totalContentPlaySessions,
97- CommonUtil .roundDouble(totalTimeSpent, 2 ), totalInteractions, totalPageviews, new DateTime ().getMillis)))
65+ if (data.count() > 0 ) {
66+ val eventDate = data.first().syncts
67+ val totalContentPlaySessions = data.filter(x =>
68+ x.dimensions.mode.getOrElse(" " ).equalsIgnoreCase(_constant.PLAY ) &&
69+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.CONTENT )).map { f =>
70+ val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
71+ eksMap.getOrElse(" total_sessions" , 0 ).asInstanceOf [Number ].longValue()
72+ }.sum().toLong
73+
74+ val totalTimeSpent = data.filter(x =>
75+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
76+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
77+ val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
78+ eksMap.getOrElse(" total_ts" , 0.0 ).asInstanceOf [Number ].doubleValue()
79+ }.sum()
80+
81+ val totalInteractions = data.filter(x =>
82+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
83+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
84+ val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
85+ eksMap.getOrElse(" total_interactions" , 0 ).asInstanceOf [Number ].longValue()
86+ }.sum().toLong
87+
88+ val totalPageviews = data.filter(x =>
89+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.APP ) ||
90+ x.dimensions.`type`.getOrElse(" " ).equalsIgnoreCase(_constant.SESSION )).map { f =>
91+ val eksMap = f.edata.eks.asInstanceOf [Map [String , AnyRef ]]
92+ eksMap.getOrElse(" total_pageviews_count" , 0 ).asInstanceOf [Number ].longValue()
93+ }.sum().toLong
94+
95+ sc.parallelize(Array (WorkFlowUsageMetricsAlgoOutput (new Date (eventDate), totalContentPlaySessions,
96+ CommonUtil .roundDouble(totalTimeSpent, 2 ), totalInteractions, totalPageviews, new DateTime ().getMillis)))
97+ }
98+ else sc.emptyRDD[WorkFlowUsageMetricsAlgoOutput ];
9899 }
99100
100101 /**
@@ -106,7 +107,7 @@ object UpdateWorkFlowUsageMetricsModel extends IBatchModelTemplate[DerivedEvent,
106107 */
107108 override def postProcess (data : RDD [WorkFlowUsageMetricsAlgoOutput ], config : Map [String , AnyRef ])
108109 (implicit sc : SparkContext ): RDD [WorkFlowUsageMetricsAlgoOutput ] = {
109- data.saveToCassandra(Constants .PLATFORM_KEY_SPACE_NAME , Constants .WORKFLOW_USAGE_SUMMARY )
110+ if (data.count() > 0 ) data.saveToCassandra(Constants .PLATFORM_KEY_SPACE_NAME , Constants .WORKFLOW_USAGE_SUMMARY )
110111 data
111112 }
112113}
0 commit comments