Similar to an above thread - I'm trying to convert...
# troubleshooting
j
Similar to an above thread - I'm trying to convert parquet to a segment. The data is in S3. I'm using Amazon EMR 6.4, which has Spark 3.1.2 (note not OSS Spark, but an Amazon fork), and EMR runs Java 8. I built Pinot 0.9.3 from source to target Java 8
mvn install package -DskipTests -Pbin-dist -Djdk.version=8
. After building, I took the contents of
pinot-distribution/target/apache-pinot-0.9.3-bin.tar.gz
and put them on HDFS (enabling all nodes to have access to the jars). I've been following this doc for general guidance. Running
spark-submit
results in
Exception in thread "main" java.lang.NoSuchMethodException: org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand.main
. Any advice? Is there a way to validate that my build is valid? My current thought is that it's either a bad build or I need to push the jars to each node and reference them locally instead of through HDFS.
m
@Xiang Fu ^^
j
Some extra info: Output from build:
Copy code
[INFO] Reactor Summary for Pinot 0.9.3:
[INFO] 
[INFO] Pinot .............................................. SUCCESS [ 20.972 s]
[INFO] Pinot Service Provider Interface ................... SUCCESS [ 19.592 s]
[INFO] Pinot Segment Service Provider Interface ........... SUCCESS [  9.869 s]
[INFO] Pinot Plugins ...................................... SUCCESS [  1.519 s]
[INFO] Pinot Metrics ...................................... SUCCESS [  1.254 s]
[INFO] Pinot Yammer Metrics ............................... SUCCESS [ 14.500 s]
[INFO] Pinot Common ....................................... SUCCESS [01:10 min]
[INFO] Pinot Input Format ................................. SUCCESS [  1.341 s]
[INFO] Pinot Avro Base .................................... SUCCESS [  8.411 s]
[INFO] Pinot Avro ......................................... SUCCESS [  5.194 s]
[INFO] Pinot Csv .......................................... SUCCESS [  4.689 s]
[INFO] Pinot JSON ......................................... SUCCESS [  4.858 s]
[INFO] Pinot local segment implementations ................ SUCCESS [ 36.507 s]
[INFO] Pinot Core ......................................... SUCCESS [ 24.983 s]
[INFO] Pinot Server ....................................... SUCCESS [ 11.239 s]
[INFO] Pinot Segment Uploader ............................. SUCCESS [  3.812 s]
[INFO] Pinot Segment Uploader Default ..................... SUCCESS [  6.803 s]
[INFO] Pinot Controller ................................... SUCCESS [01:28 min]
[INFO] Pinot Broker ....................................... SUCCESS [ 11.795 s]
[INFO] Pinot Clients ...................................... SUCCESS [  0.795 s]
[INFO] Pinot Java Client .................................. SUCCESS [  6.089 s]
[INFO] Pinot JDBC Client .................................. SUCCESS [  5.431 s]
[INFO] Pinot Batch Ingestion .............................. SUCCESS [  3.128 s]
[INFO] Pinot Batch Ingestion Common ....................... SUCCESS [  5.587 s]
[INFO] Pinot Minion ....................................... SUCCESS [  4.491 s]
[INFO] Pinot Confluent Avro ............................... SUCCESS [  4.272 s]
[INFO] Pinot ORC .......................................... SUCCESS [ 18.907 s]
[INFO] Pinot Parquet ...................................... SUCCESS [ 18.637 s]
[INFO] Pinot Thrift ....................................... SUCCESS [  4.442 s]
[INFO] Pinot Protocol Buffers ............................. SUCCESS [  4.660 s]
[INFO] Pluggable Pinot file system ........................ SUCCESS [  0.636 s]
[INFO] Pinot Azure Data Lake Storage ...................... SUCCESS [ 11.405 s]
[INFO] Pinot Hadoop Filesystem ............................ SUCCESS [  5.226 s]
[INFO] Pinot Google Cloud Storage ......................... SUCCESS [  8.767 s]
[INFO] Pinot Amazon S3 .................................... SUCCESS [ 14.118 s]
[INFO] Pinot Batch Ingestion for Spark .................... SUCCESS [ 12.647 s]
[INFO] Pinot Batch Ingestion for Hadoop ................... SUCCESS [ 11.989 s]
[INFO] Pinot Batch Ingestion Standalone ................... SUCCESS [  6.444 s]
[INFO] Pinot Batch Ingestion .............................. SUCCESS [  3.006 s]
[INFO] Pinot Ingestion Common ............................. SUCCESS [ 10.756 s]
[INFO] Pinot Hadoop ....................................... SUCCESS [ 43.413 s]
[INFO] Pinot Spark ........................................ SUCCESS [ 49.759 s]
[INFO] Pinot Stream Ingestion ............................. SUCCESS [  0.664 s]
[INFO] Pinot Kafka Base ................................... SUCCESS [  3.071 s]
[INFO] Pinot Kafka 0.9 .................................... SUCCESS [ 11.378 s]
[INFO] Pinot Kafka 2.0 .................................... SUCCESS [ 11.595 s]
[INFO] Pinot Kinesis ...................................... SUCCESS [  9.517 s]
[INFO] Pinot Pulsar ....................................... SUCCESS [ 27.551 s]
[INFO] Pinot Minion Tasks ................................. SUCCESS [  4.847 s]
[INFO] Pinot Minion Built-In Tasks ........................ SUCCESS [  8.724 s]
[INFO] Pinot Dropwizard Metrics ........................... SUCCESS [  8.238 s]
[INFO] Pinot Segment Writer ............................... SUCCESS [  4.994 s]
[INFO] Pinot Segment Writer File Based .................... SUCCESS [  7.465 s]
[INFO] Pluggable Pinot Environment Provider ............... SUCCESS [  0.681 s]
[INFO] Pinot Azure Environment ............................ SUCCESS [ 29.145 s]
[INFO] Pinot Tools ........................................ SUCCESS [01:40 min]
[INFO] Pinot Test Utils ................................... SUCCESS [ 17.145 s]
[INFO] Pinot Integration Tests ............................ SUCCESS [ 19.640 s]
[INFO] Pinot Perf ......................................... SUCCESS [01:03 min]
[INFO] Pinot Distribution ................................. SUCCESS [01:45 min]
[INFO] Pinot Connectors ................................... SUCCESS [  0.333 s]
[INFO] Pinot Spark Connector .............................. SUCCESS [ 31.685 s]
[INFO] Presto Pinot Driver ................................ SUCCESS [  4.653 s]
[INFO] ------------------------------------------------------------------------
[INFO] BUILD SUCCESS
Spark submit w/ stack trace:
Copy code
sh-4.2$ export PINOT_DISTRIBUTION_DIR=hdfs:///apps/pinot/apache-pinot-${PINOT_VERSION}-bin
sh-4.2$ spark-submit \
>     --class org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand \
>     --master local --deploy-mode client \
>     --conf "spark.driver.extraJavaOptions=-Dplugins.dir=${PINOT_DISTRIBUTION_DIR}/plugins -Dplugins.include=pinot-s3,pinot-parquet -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml" \
>     --conf "spark.driver.extraClassPath=${PINOT_DISTRIBUTION_DIR}/plugins/pinot-batch-ingestion/pinot-batch-ingestion-spark/pinot-batch-ingestion-spark-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-file-system/pinot-s3/pinot-s3-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-input-format/pinot-parquet/pinot-parquet-${PINOT_VERSION}-shaded.jar" \
>     ${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar -jobSpecFile job_spec.yaml
SLF4J: Class path contains multiple SLF4J bindings.
SLF4J: Found binding in [jar:file:/usr/lib/spark/jars/slf4j-log4j12-1.7.30.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/usr/share/aws/emr/emrfs/lib/slf4j-log4j12-1.7.12.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/usr/share/aws/redshift/jdbc/redshift-jdbc42-1.2.37.1061.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: See <http://www.slf4j.org/codes.html#multiple_bindings> for an explanation.
SLF4J: Actual binding is of type [org.slf4j.impl.Log4jLoggerFactory]
Exception in thread "main" java.lang.NoSuchMethodException: org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand.main([Ljava.lang.String;)
	at java.lang.Class.getMethod(Class.java:1814)
	at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:42)
	at <http://org.apache.spark.deploy.SparkSubmit.org|org.apache.spark.deploy.SparkSubmit.org>$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:959)
	at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:180)
	at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:203)
	at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:90)
	at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1047)
	at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1056)
	at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
22/02/10 22:06:51 INFO ShutdownHookManager: Shutdown hook called
22/02/10 22:06:51 INFO ShutdownHookManager: Deleting directory /mnt/tmp/spark-0ceac8e4-b275-4047-953e-9fe0219866cb
Job spec:
Copy code
executionFrameworkSpec:
  name: 'spark'
  segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentGenerationJobRunner'
jobType: SegmentCreation
extraConfigs:
    stagingDir: '...'
inputDirURI: '...'
outputDirURI: '...'
includeFileNamePattern: 'glob:**/*.parquet'
overwriteOutput: true
pinotFSSpecs:
  - scheme: s3
    className: org.apache.pinot.plugin.filesystem.S3PinotFS
    configs:    
      region: 'us-east-1'
recordReaderSpec:
  dataFormat: 'parquet'
  className: 'org.apache.pinot.plugin.inputformat.parquet.ParquetRecordReader'
segmentNameGeneratorSpec:
  type: normalizedDate
I've verified that the HDFS path is valid and has the expected contents. I also opened up the jar file and verified that the file that it says is missing does in fact existing.
x
This was due to that the main method is removed
šŸ‘ 1
I just added that back yesterday
can you either try 0.8.0
or We will have it back in 0.10.0
j
Are there any issues with creating the segments in 0.8, but having the server be 0.9?
x
no, just some 0.9 features are not there, which should be minimal
šŸ‘ 1
m
You should be using
PinotAdministrator
class with passing LaunchDataIngestionJob as arg.
j
I tried 0.9 w/ the above suggestion and it seems to fail to instantiate the class:
Copy code
spark-submit \
    --class org.apache.pinot.tools.admin.PinotAdministrator \
    --master local --deploy-mode client \
    --conf "spark.driver.extraJavaOptions=-Dplugins.dir=${PINOT_DISTRIBUTION_DIR}/plugins -Dplugins.include=pinot-s3,pinot-parquet -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml" \
    --conf "spark.driver.extraClassPath=${PINOT_DISTRIBUTION_DIR}/plugins/pinot-batch-ingestion/pinot-batch-ingestion-spark/pinot-batch-ingestion-spark-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-file-system/pinot-s3/pinot-s3-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-input-format/pinot-parquet/pinot-parquet-${PINOT_VERSION}-shaded.jar" \
    ${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar LaunchDataIngestionJob -jobSpecFile job_spec.yaml
Copy code
Exception in thread "main" java.lang.ExceptionInInitializerError
	at org.apache.pinot.tools.admin.command.StartKafkaCommand.<init>(StartKafkaCommand.java:51)
	at org.apache.pinot.tools.admin.PinotAdministrator.<clinit>(PinotAdministrator.java:98)
	at java.lang.Class.forName0(Native Method)
	at java.lang.Class.forName(Class.java:348)
	at org.apache.spark.util.Utils$.classForName(Utils.scala:207)
	at <http://org.apache.spark.deploy.SparkSubmit.org|org.apache.spark.deploy.SparkSubmit.org>$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:924)
	at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:180)
	at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:203)
	at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:90)
	at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1047)
	at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1056)
	at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
Caused by: java.util.NoSuchElementException
	at java.util.ServiceLoader$LazyIterator.nextService(ServiceLoader.java:365)
	at java.util.ServiceLoader$LazyIterator.next(ServiceLoader.java:404)
	at java.util.ServiceLoader$1.next(ServiceLoader.java:480)
	at org.apache.pinot.tools.utils.KafkaStarterUtils.getKafkaConnectorPackageName(KafkaStarterUtils.java:54)
	at org.apache.pinot.tools.utils.KafkaStarterUtils.<clinit>(KafkaStarterUtils.java:46)
	... 12 more
I'll likely poke at this a bit more and in parallel build 0.8
Using 0.8 seems to work šŸŽ‰ Once I successfully get a segment produced I'll post the steps here to help anyone else who stumbles upon this thread. It's a bit different from what's in the doc.
x
Thanks! The Main method is accidentally deleted in 0.9. We are adding it back in 0.10.0
šŸ‘ 1
j
Update: • The job doesn't seem to work with yarn • Running in local mode also fails (potentially incompatible with Spark 3 or Amazon's fork, not sure) • The plugin manager doesn't seem to support files located in HDFS, i.e. when
-Dplugins.dir
points to a location in HDFS an incorrect warning stating the directory doesn't exist occurs. This is worked around by having the jars locally on each node (In EMR this would necessitate using a bootstrap script) I'm trying with an older version of EMR which runs Spark 2: emr-5.34.0 and Spark 2.4.8.
With emr-5.34.0 the issue with YARN is gone. I'm not hitting the same issue I ran into when running locally (same error with YARN or scheduled locally). I'm assuming this is a result of Amazon having a fork:
INFO [SparkContext] [main] Running Spark version 2.4.8-amzn-0
.
Copy code
Exception in thread "main" java.lang.VerifyError: Bad return type
Exception Details:
  Location:
    org/apache/hadoop/hdfs/DFSClient.getQuotaUsage(Ljava/lang/String;)Lorg/apache/hadoop/fs/QuotaUsage; @157: areturn
  Reason:
    Type 'org/apache/hadoop/fs/ContentSummary' (current frame, stack[0]) is not assignable to 'org/apache/hadoop/fs/QuotaUsage' (from method signature)
  Current Frame:
    bci: @157
    flags: { }
    locals: { 'org/apache/hadoop/hdfs/DFSClient', 'java/lang/String', 'org/apache/hadoop/ipc/RemoteException', 'java/io/IOException' }
    stack: { 'org/apache/hadoop/fs/ContentSummary' }
  Bytecode:
    0x0000000: 2ab6 00de 2a13 023b 2bb6 00b3 4d01 4e2a
    0x0000010: b400 432b b902 3c02 003a 042c c600 1d2d
    0x0000020: c600 152c b600 b5a7 0012 3a05 2d19 05b6
    0x0000030: 00b7 a700 072c b600 b519 04b0 3a04 1904
    0x0000040: 4e19 04bf 3a06 2cc6 001d 2dc6 0015 2cb6
    0x0000050: 00b5 a700 123a 072d 1907 b600 b7a7 0007
    0x0000060: 2cb6 00b5 1906 bf4d 2c07 bd00 d059 0312
    0x0000070: d253 5904 12dc 5359 0512 dd53 5906 1302
    0x0000080: 3d53 b600 d34e 2dc1 023d 9900 14b2 0024
    0x0000090: 1302 3eb9 002c 0200 2a2b b602 3fb0 2dbf
    0x00000a0:                                        
  Exception Handler Table:
    bci [35, 39] => handler: 42
    bci [15, 27] => handler: 60
    bci [15, 27] => handler: 68
    bci [78, 82] => handler: 85
    bci [60, 70] => handler: 68
    bci [4, 57] => handler: 103
    bci [60, 103] => handler: 103
  Stackmap Table:
    full_frame(@42,{Object[#818],Object[#839],Object[#893],Object[#864],Object[#1335]},{Object[#864]})
    same_frame(@53)
    same_frame(@57)
    full_frame(@60,{Object[#818],Object[#839],Object[#893],Object[#864]},{Object[#864]})
    same_locals_1_stack_item_frame(@68,Object[#864])
    full_frame(@85,{Object[#818],Object[#839],Object[#893],Object[#864],Top,Top,Object[#864]},{Object[#864]})
    same_frame(@96)
    same_frame(@100)
    full_frame(@103,{Object[#818],Object[#839]},{Object[#918]})
    append_frame(@158,Object[#918],Object[#879])

	at org.apache.hadoop.hdfs.DistributedFileSystem.initialize(DistributedFileSystem.java:159)
	at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2653)
	at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:92)
	at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2687)
	at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2669)
	at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:371)
	at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:362)
	at org.apache.spark.util.Utils$.getHadoopFileSystem(Utils.scala:1911)
	at org.apache.spark.deploy.history.EventLogFileWriter.<init>(EventLogFileWriters.scala:60)
	at org.apache.spark.deploy.history.SingleEventLogFileWriter.<init>(EventLogFileWriters.scala:211)
	at org.apache.spark.deploy.history.EventLogFileWriter$.apply(EventLogFileWriters.scala:181)
	at org.apache.spark.scheduler.EventLoggingListener.<init>(EventLoggingListener.scala:73)
	at org.apache.spark.SparkContext.<init>(SparkContext.scala:535)
	at org.apache.spark.SparkContext.<init>(SparkContext.scala:118)
	at org.apache.spark.SparkContext$.getOrCreate(SparkContext.scala:2595)
	at org.apache.spark.SparkContext.getOrCreate(SparkContext.scala)
	at org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentGenerationJobRunner.run(SparkSegmentGenerationJobRunner.java:196)
	at org.apache.pinot.spi.ingestion.batch.IngestionJobLauncher.kickoffIngestionJob(IngestionJobLauncher.java:142)
	at org.apache.pinot.spi.ingestion.batch.IngestionJobLauncher.runIngestionJob(IngestionJobLauncher.java:101)
	at org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand.execute(LaunchDataIngestionJobCommand.java:132)
	at org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand.main(LaunchDataIngestionJobCommand.java:67)
	at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
	at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.lang.reflect.Method.invoke(Method.java:498)
	at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52)
	at <http://org.apache.spark.deploy.SparkSubmit.org|org.apache.spark.deploy.SparkSubmit.org>$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:863)
	at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:161)
	at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:184)
	at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:86)
	at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:938)
	at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:947)
	at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
Assuming this is the case, then I think it requires deploying OSS Spark 2 across the cluster and not using Amazon's version. I ran into this same problem before (Amazon's version conflicting with a different OSS project).
I gave up on spark for now and instead I'm doing the conversion using the standalone tool. It's not ideal, but can be made better by doing the following: 1. Set
segmentCreationJobParallelism
in the job spec (parallelism seems to be limited by the number of parquet files). In my setup this can do ~250k records / sec on a single core. 2. Set
java.io.tmpdir
to a location with enough space (in my case a separately mounted EBS volume). Example:
Copy code
export JAVA_OPTS="-Xms4G -Dlog4j2.configurationFile=conf/log4j2.xml -Dpinot.admin.system.exit=true -Djava.io.tmpdir=/data/scratch"
pinot-admin.sh LaunchDataIngestionJob -jobSpecFile /data/pinot-config/job_spec.yaml
x
Thanks @James Mnatzaganian, pinot ingestion story is a bit complicated with spark as we actually need to package all the required jars together into a uber jar. E.g. if you input is orc files, then you need to package pinot-orc plugin, but that also brings hadoop dependency transparently. Hence we choose the plugin directory upload model, but seems it also has a lot of issue.
šŸ‘ 1
cc: @Dunith Dhanushka
j
Dependency management for the hadoop ecosystem is tricky. From what I've seen if all the versions don't align perfectly then it doesn't work. This is further complicated by vendors having their own forks. If there's a relatively low friction path to make it such that one could build the job to match the necessary versions of the deployment environment that could solve some of these issues. From a production standpoint, if I were to create segments offline, it would definitely be done via Spark. For someone already established with Spark what was expected to be a low-friction path ended up becoming high friction -- I would now need to have a new environment just for building segments. --- Looking at the source code -- it seems relatively straightforward in that for each file a segment is created. From what I can tell, all that's really needed is to be able to take the input files and map them to a task spec. My personal preference would be to instead take a DataFrame and have Pinot extend the functionality by taking a config and producing the dataset as needed. This might be a bit strange (unless reading segments and loading them as DFs is also supported, which could be a nice way to easily take data out of Pinot), but it'd also enable you to go beyond 1 parquet file to 1 segment and instead have segments built in a way that is known to be most optimal for Pinot with the relevant user parameters exposed. It also means that as a Spark user I can trivially add whatever transformations I want in the same job and performance wise it doesn't matter if I'm using Java/Scala/Python since the heavy lifting is covered by the underlying Pinot converter. Another benefit is that you don't have to worry about the input data source, since you'd directly be given the DF.
m
I think LinkedIn did this @Jack ?
j
Right, LinkedIn does maintain an internal code that takes the dataFrame instead of an input path as the input and then generates the Pinot segment from the DF. The core logic is that inside each of the partitions of DF, run the similar logic like this: https://github.com/apache/pinot/blob/08b909c45e85f9bf8d8659561a2d13b4cc443ebc/pino[…]segment/local/indexsegment/mutable/IntermediateSegmentTest.java
j
Thanks! HW for later, assuming the POC progress beyond a POC šŸ™‚ I appreciate the tip!
šŸ‘ 2