Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 25 additions & 5 deletions dfutil/enrolment/acbp/acbpDFUtil.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,14 +13,25 @@
sys.path.append(str(Path(__file__).resolve().parents[3]))
from util import schemas
from constants.ParquetFileConstants import ParquetFileConstants
from pyspark.sql import functions as F


def preComputeACBPData(spark):
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false")
spark.conf.set("spark.sql.parquet.outputTimestampType", "TIMESTAMP_MICROS")
acbp_df = spark.read.parquet(ParquetFileConstants.ACBP_PARQUET_FILE)

acbp_select_df = acbp_df \
.withColumn("orgid", explode(col("orgidlist"))) \
.withColumn("context_data", from_json(col("contextdata"), schemas.accessControlSchema)) \
.withColumn("userGroup", explode(col("context_data.accessControl.userGroups"))) \
.withColumn("criteria", explode(col("userGroup.userGroupCriteriaList"))) \
.withColumn("assignmentType",
when(col("criteria.criteriaKey") == "rootOrgId", "AllUser")
.when(col("criteria.criteriaKey") == "designation", "Designation")
.when(col("criteria.criteriaKey") == "user", "CustomUser")
.otherwise(col("criteria.criteriaKey"))
) \
.select(
col("planid").alias("acbpID"),
col("orgid").alias("userOrgID"),
Expand All @@ -32,12 +43,13 @@ def preComputeACBPData(spark):
col("enddate").alias("completionDueDate"),
col("publishedat").alias("allocatedOn"),
col("contentlist").alias("acbpCourseIDList"),
lit(None).alias("assignmentType"),
lit(None).alias("assignmentTypeInfo")
col("assignmentType"),
col("criteria.criteriaValue").alias("assignmentTypeInfo")
) \
.na.fill({"cbPlanName": ""})

# Draft data
draft_cbp_data = acbp_select_df.filter((col("acbpStatus") == "DRAFT") & col("draftdata").isNotNull()) \
draft_cbp_data = acbp_select_df.filter((col("acbpStatus") == "draft") & col("draftdata").isNotNull()) \
.select("acbpID", "userOrgID", "draftdata", "acbpStatus", "acbpCreatedBy", "isapar") \
.withColumn("draftData", from_json(col("draftdata"), schemas.cbplan_draft_data_schema)) \
.withColumn("cbPlanName", col("draftData.name")) \
Expand All @@ -47,11 +59,15 @@ def preComputeACBPData(spark):
.withColumn("allocatedOn", lit("not published")) \
.withColumn("acbpCourseIDList", col("draftData.contentList")) \
.drop("draftData")

# Non-draft data
non_draft_cbp_data = acbp_select_df.filter(col("acbpStatus") != "DRAFT")
non_draft_cbp_data = acbp_select_df.filter(col("acbpStatus") != "draft")

# Add draftdata column with null values to non-draft data
draft_cbp_data = draft_cbp_data.withColumn("draftdata", lit(None).cast("string"))

# Union the two
final_df = non_draft_cbp_data.unionByName(draft_cbp_data)
final_df = non_draft_cbp_data.unionByName(draft_cbp_data)

exportDFToParquet(final_df, ParquetFileConstants.ACBP_SELECT_FILE)
explodeAcbpData(spark, final_df)
Expand All @@ -64,6 +80,8 @@ def explodeAcbpData(spark, acbp_df):
"assignmentType", "completionDueDate", "allocatedOn", "acbpCourseIDList", "acbpStatus",
"acbpCreatedBy", "cbPlanName"]
user_df = spark.read.parquet(ParquetFileConstants.USER_ORG_COMPUTED_FILE)

# CustomUser
acbp_custom_user_df = acbp_df \
.filter(col("assignmentType") == "CustomUser") \
.withColumn("userID", explode(col("assignmentTypeInfo"))) \
Expand All @@ -74,10 +92,12 @@ def explodeAcbpData(spark, acbp_df):
.filter(col("assignmentType") == "Designation") \
.withColumn("designation", explode(col("assignmentTypeInfo"))) \
.join(user_df, ["userOrgID", "designation"], "left")

# AllUser
acbp_all_user_df = acbp_df \
.filter(col("assignmentType") == "AllUser") \
.join(user_df, ["userOrgID"], "left")

# Union the three DataFrames
dfs = [acbp_custom_user_df, acbp_designation_df, acbp_all_user_df]
selected_dfs = [df.select([col(c) for c in selectColumns]) for df in dfs]
Expand Down
5 changes: 4 additions & 1 deletion jobs/stage-2/l2Assessments.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@ def process_report(self,spark,config):
contentDF.show(5, truncate=False)

# user details dataframe
#userDF = userDF.filter((col("status") == 1) & (col("mdo_id") == '0135502316148080641003'))
userDF = userDF.filter(col("status") == 1)

print("userDF count:", userDF.count())
Expand Down Expand Up @@ -276,7 +277,9 @@ def process_report(self,spark,config):
print("Stage 10: Generating final report...")
# Export report
apar_assessment_data.coalesce(1).write.mode("overwrite").parquet("/mount/data/analytics/igot-reports/assessment-report-apar/parquet")

#csv
#apar_assessment_data.coalesce(1).write.mode("overwrite").option("header", "true").csv("/home/analytics/shishir/assessment-report-apar/csv")

except Exception as e:
print(f"Error occurred during processing: {str(e)}")
raise
Expand Down
Loading