Question on importing partitioned parquet data fro...
# general
s
Question on importing partitioned parquet data from S3… I have looked at the docs, examples, and searched Slack, but did not find an answer. I have a synthetic data set based on the transcript data example from the docs. That data is produced in Spark.
Copy code
transcripts_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
Copy code
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
Copy code
{
  "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.
Copy code
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 subject
Re: Importing partitioned data from S3
Copy code
This 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
Copy code
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…
Data in the parquet file shows all columns
S3 View
a
fwiw #C011C9JHN7R might be a bit better of a channel, but what's your table setup like?
namely, what's your time column set there?
s
Hi Ashwin. Table config
Copy code
{
  "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
  }
}
a
and what's your job spec file?
s
transcript_job_spec.yml
a
@Steven Hall one thing is that you dont have
stagingDir
set?
which you need when generating segments in spark
s
yes
a
sorry, yes you have it set somewhere else, or yes you still need to set it?
s
OK, tell me more. It seems to work, until partitioned in a specific manner. I can import 250K records, until I change the partitioning. Let me look at staging dir. Do you have a doc link I should look at?
search stagingDir on there
We also need a temporary
stagingDir
for our spark job
i 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 limitiation
(i wouldnt recommend it tho, you really want to partition by time range if possible for Pinot)
s
I think your last comment likely gets to the root of the issue. I have been thinking about this wrong. It should be partitioned by time for Pinot. Thanks, let me investigate that feedback a bit. Appreciate the help.
a
yeah, and specifically by time range is ideal
since then Pinot can prune segments at query time
you can ofc partition by other stuff, but generally folks do time partitioning, and that'll be the default setup when using a time column like that
s
👍