Steven Hall
01/27/2023, 5:47 PMtranscripts_bucket = "<s3://transcript-parquet/1/>"
transcript_df.write.mode("overwrite").parquet(transcripts_bucket)
I run a Spark job to transform into segments and import the data into Pinot
spark_args = {
'master': '<spark://spark-master:7077>',
'deploy_mode': 'cluster',
'name': 'segments-from-parquet',
'class': 'org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand',
'executor_memory': '2G',
'executor_cores': '1',
'total_executor_cores': '2',
'verbose': True,
'conf': [f"spark.driver.extraJavaOptions='{EXTRA_JAVA_OPTIONS}'"],
'main_file_args': '-jobSpecFile=/home/job_specs/transcript_job_spec.yml'
}
main_file = f'{PINOT_DISTRIBUTION_DIR}/lib/pinot-all-{PINOT_VERSION}-jar-with-dependencies.jar'
app = SparkJob(main_file, **spark_args)
app.submit()
All columns populate as expected. Schema is
{
"schemaName": "transcript_indexed",
"dimensionFieldSpecs": [
{
"name": "studentID",
"dataType": "INT"
},
{
"name": "firstName",
"dataType": "STRING"
},
{
"name": "lastName",
"dataType": "STRING"
},
{
"name": "gender",
"dataType": "STRING"
},
{
"name": "subject",
"dataType": "STRING"
}
],
"metricFieldSpecs": [
{
"name": "score",
"dataType": "FLOAT"
}
],
"timeFieldSpec": {
"incomingGranularitySpec": {
"name": "examTime",
"dataType": "LONG",
"timeType": "MILLISECONDS"
}
}
}
If I partition my data by subject when I write the parquet files, I have an unexpected outcome, the subject field in the Pinot segments is null.
transcripts_bucket = "<s3://transcript-parquet/1/>"
transcript_df.write.mode("overwrite").partitionBy("subject").parquet(transcripts_bucket)
Are we thinking about this incorrectly…. in a way that Pinot does not support? Alternately, is there some change we need to make in the configs to work with data lake data that is normally partitioned?
The data on Minio — my S3 service fake looks like this once partitioned by subjectSteven Hall
01/27/2023, 6:36 PMThis is getting even more interesting. If I change the partitioning in Spark as follows
transcripts_bucket = "<s3://transcript-parquet/1/>"
transcript_df.repartition("subject").write.mode("overwrite").parquet(transcripts_bucket)
The data does not import. Error in the Spark worker is
Caused by: java.lang.NullPointerException: columns should not be null
at org.apache.pinot.shaded.org.apache.parquet.Preconditions.checkNotNull(Preconditions.java:36)
S3 now looks like this…Steven Hall
01/27/2023, 6:40 PMSteven Hall
01/27/2023, 6:44 PMAshwin Raja
01/27/2023, 7:05 PMAshwin Raja
01/27/2023, 7:05 PMSteven Hall
01/27/2023, 7:09 PM{
"OFFLINE": {
"tableName": "transcript_indexed_OFFLINE",
"tableType": "OFFLINE",
"segmentsConfig": {
"replication": "1",
"timeType": "MILLISECONDS",
"schemaName": "transcript_indexed",
"timeColumnName": "examTime",
"minimizeDataMovement": false
},
"tenants": {
"broker": "DefaultTenant",
"server": "DefaultTenant"
},
"tableIndexConfig": {
"invertedIndexColumns": [],
"rangeIndexVersion": 2,
"loadMode": "MMAP",
"enableDefaultStarTree": false,
"aggregateMetrics": false,
"nullHandlingEnabled": false,
"autoGeneratedInvertedIndex": false,
"enableDynamicStarTreeCreation": false,
"optimizeDictionaryForMetrics": false,
"noDictionarySizeRatioThreshold": 0,
"createInvertedIndexDuringSegmentGeneration": false
},
"metadata": {},
"isDimTable": false
}
}Ashwin Raja
01/27/2023, 7:14 PMSteven Hall
01/27/2023, 7:17 PMAshwin Raja
01/27/2023, 7:19 PMstagingDir set?Ashwin Raja
01/27/2023, 7:20 PMSteven Hall
01/27/2023, 7:20 PMAshwin Raja
01/27/2023, 7:22 PMSteven Hall
01/27/2023, 7:23 PMAshwin Raja
01/27/2023, 7:23 PMAshwin Raja
01/27/2023, 7:23 PMAshwin Raja
01/27/2023, 7:24 PMWe also need a temporaryi dont really know if related, but we were able to create segments on a non-time-partitioned table in spark just fine, so I don't think it's a hard limitiationfor our spark jobstagingDir
Ashwin Raja
01/27/2023, 7:24 PMSteven Hall
01/27/2023, 7:26 PMAshwin Raja
01/27/2023, 7:26 PMAshwin Raja
01/27/2023, 7:26 PMAshwin Raja
01/27/2023, 7:27 PMSteven Hall
01/27/2023, 7:27 PM