soumilshah1995 opened a new issue, #10177:
URL: https://github.com/apache/hudi/issues/10177

   Dear Community,
   
   I hope this message finds you well. I am currently facing an issue with the 
Delta Streamer, particularly around the schema registry. The error message I'm 
encountering is:
   
   ```
   org.apache.hudi.internal.schema.HoodieSchemaException: Failed to parse 
schema from registry:
   
   ```
   
   Here is a brief overview of my stack:
   
   Docker for Kafka, ZooKeeper, Debezium, Schema Registry, and Postgres.
   Debezium connector configuration:
   ```
   name=PostgresConnector
   connector.class=io.debezium.connector.postgresql.PostgresConnector
   database.user=postgres
   database.dbname=postgres
   tasks.max=1
   database.hostname=postgres
   database.password=postgres
   database.server.name=postgres
   table.include.list=public.sales
   database.port=5432
   plugin.name=pgoutput
   ```
   
   I have successfully created a table in Postgres, inserted records, and 
verified messages in the UI. However, the issue arises when running the Delta 
Streamer job using the following command:
   
   <img width="1204" alt="Screenshot 2023-11-26 at 9 32 21 AM" 
src="https://github.com/apache/hudi/assets/39345855/57037c31-19de-4ceb-a626-718be01f7abc";>
   
   
   ```
   
   # Delta Streamer Job
   spark-submit \
       --class org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer \
       --packages org.apache.hudi:hudi-spark3.4-bundle_2.12:0.14.0 \
       --properties-file spark-config.properties \
       --master 'local[*]' \
       --executor-memory 1g \
       jar/hudi-utilities-slim-bundle_2.12-0.14.0.jar \
       --table-type COPY_ON_WRITE \
       --op UPSERT \
       --source-ordering-field ts \
       --source-class org.apache.hudi.utilities.sources.AvroKafkaSource \
       --target-base-path 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
 \
       --target-table orders \
       --schemaprovider-class 
org.apache.hudi.utilities.schema.SchemaRegistryProvider \
       --props hudi_tbl.props
   
   ```
   #### Note 
   
   Tried using this as well in spark submit 
   ```
       --source-class 
org.apache.hudi.utilities.sources.debezium.PostgresDebeziumSource \
       --payload-class 
org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload \
       
   ```
   
   These are my Hudi settings related to the schema registry:
   
   
   ```
   
   hoodie.datasource.write.recordkey.field=order_id
   hoodie.datasource.write.partitionpath.field=order_date
   hoodie.datasource.write.precombine.field=ts
   bootstrap.servers=localhost:7092
   auto.offset.reset=earliest
   hoodie.deltastreamer.source.kafka.topic=postgres.public.sales
   
   
hoodie.deltastreamer.source.kafka.value.deserializer.class=org.apache.hudi.utilities.deser.KafkaAvroSchemaDeserializer
   schema.registry.url=http://localhost:8081
   
hoodie.deltastreamer.schemaprovider.registry.url=http://localhost:8081/subjects/postgres.public.sales-value/versions/latest
   
   ```
   
   The error I'm encountering specifically mentions an issue with fetching the 
schema from the registry:
   
   Error reading source schema from registry : 
http://localhost:8081/subjects/postgres.public.sales-value/versions/latest
   
   # Error Logs 
   ```
   
   inuing without scheduling configs
   23/11/26 09:24:53 INFO SparkContext: Running Spark version 3.4.0
   23/11/26 09:24:53 INFO ResourceUtils: 
==============================================================
   23/11/26 09:24:53 INFO ResourceUtils: No custom resources configured for 
spark.driver.
   23/11/26 09:24:53 INFO ResourceUtils: 
==============================================================
   23/11/26 09:24:53 INFO SparkContext: Submitted application: streamer-orders
   23/11/26 09:24:54 INFO ResourceProfile: Default ResourceProfile created, 
executor resources: Map(cores -> name: cores, amount: 1, script: , vendor: , 
memory -> name: memory, amount: 1024, script: , vendor: , offHeap -> name: 
offHeap, amount: 0, script: , vendor: ), task resources: Map(cpus -> name: 
cpus, amount: 1.0)
   23/11/26 09:24:54 INFO ResourceProfile: Limiting resource is cpu
   23/11/26 09:24:54 INFO ResourceProfileManager: Added ResourceProfile id: 0
   23/11/26 09:24:54 INFO SecurityManager: Changing view acls to: soumilshah
   23/11/26 09:24:54 INFO SecurityManager: Changing modify acls to: soumilshah
   23/11/26 09:24:54 INFO SecurityManager: Changing view acls groups to: 
   23/11/26 09:24:54 INFO SecurityManager: Changing modify acls groups to: 
   23/11/26 09:24:54 INFO SecurityManager: SecurityManager: authentication 
disabled; ui acls disabled; users with view permissions: soumilshah; groups 
with view permissions: EMPTY; users with modify permissions: soumilshah; groups 
with modify permissions: EMPTY
   23/11/26 09:24:54 INFO deprecation: mapred.output.compression.codec is 
deprecated. Instead, use mapreduce.output.fileoutputformat.compress.codec
   23/11/26 09:24:54 INFO deprecation: mapred.output.compress is deprecated. 
Instead, use mapreduce.output.fileoutputformat.compress
   23/11/26 09:24:54 INFO deprecation: mapred.output.compression.type is 
deprecated. Instead, use mapreduce.output.fileoutputformat.compress.type
   23/11/26 09:24:54 INFO Utils: Successfully started service 'sparkDriver' on 
port 64395.
   23/11/26 09:24:54 INFO SparkEnv: Registering MapOutputTracker
   23/11/26 09:24:54 INFO SparkEnv: Registering BlockManagerMaster
   23/11/26 09:24:54 INFO BlockManagerMasterEndpoint: Using 
org.apache.spark.storage.DefaultTopologyMapper for getting topology information
   23/11/26 09:24:54 INFO BlockManagerMasterEndpoint: 
BlockManagerMasterEndpoint up
   23/11/26 09:24:54 INFO SparkEnv: Registering BlockManagerMasterHeartbeat
   23/11/26 09:24:54 INFO DiskBlockManager: Created local directory at 
/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/blockmgr-db92f060-3d76-4395-be18-7e06f038f57e
   23/11/26 09:24:54 INFO MemoryStore: MemoryStore started with capacity 434.4 
MiB
   23/11/26 09:24:54 INFO SparkEnv: Registering OutputCommitCoordinator
   23/11/26 09:24:54 INFO JettyUtils: Start Jetty 0.0.0.0:8090 for SparkUI
   23/11/26 09:24:54 INFO Utils: Successfully started service 'SparkUI' on port 
8090.
   23/11/26 09:24:54 INFO SparkContext: Added JAR 
file:///Users/soumilshah/.ivy2/jars/org.apache.hudi_hudi-spark3.4-bundle_2.12-0.14.0.jar
 at 
spark://soumils-mbp:64395/jars/org.apache.hudi_hudi-spark3.4-bundle_2.12-0.14.0.jar
 with timestamp 1701008693979
   23/11/26 09:24:54 INFO SparkContext: Added JAR 
file:/Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/jar/hudi-utilities-slim-bundle_2.12-0.14.0.jar
 at spark://soumils-mbp:64395/jars/hudi-utilities-slim-bundle_2.12-0.14.0.jar 
with timestamp 1701008693979
   23/11/26 09:24:54 INFO Executor: Starting executor ID driver on host 
soumils-mbp
   23/11/26 09:24:54 INFO Executor: Starting executor with user classpath 
(userClassPathFirst = false): ''
   23/11/26 09:24:54 INFO Executor: Fetching 
spark://soumils-mbp:64395/jars/org.apache.hudi_hudi-spark3.4-bundle_2.12-0.14.0.jar
 with timestamp 1701008693979
   23/11/26 09:24:54 INFO TransportClientFactory: Successfully created 
connection to soumils-mbp/192.168.1.31:64395 after 14 ms (0 ms spent in 
bootstraps)
   23/11/26 09:24:54 INFO Utils: Fetching 
spark://soumils-mbp:64395/jars/org.apache.hudi_hudi-spark3.4-bundle_2.12-0.14.0.jar
 to 
/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-3f9f543b-ad9d-459b-95cc-19fc7f2573e0/userFiles-ffd39215-161f-4f21-b09f-3536e270730a/fetchFileTemp4096719352216213838.tmp
   23/11/26 09:24:54 INFO Executor: Adding 
file:/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-3f9f543b-ad9d-459b-95cc-19fc7f2573e0/userFiles-ffd39215-161f-4f21-b09f-3536e270730a/org.apache.hudi_hudi-spark3.4-bundle_2.12-0.14.0.jar
 to class loader
   23/11/26 09:24:54 INFO Executor: Fetching 
spark://soumils-mbp:64395/jars/hudi-utilities-slim-bundle_2.12-0.14.0.jar with 
timestamp 1701008693979
   23/11/26 09:24:54 INFO Utils: Fetching 
spark://soumils-mbp:64395/jars/hudi-utilities-slim-bundle_2.12-0.14.0.jar to 
/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-3f9f543b-ad9d-459b-95cc-19fc7f2573e0/userFiles-ffd39215-161f-4f21-b09f-3536e270730a/fetchFileTemp7800412205334403898.tmp
   23/11/26 09:24:54 INFO Executor: Adding 
file:/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-3f9f543b-ad9d-459b-95cc-19fc7f2573e0/userFiles-ffd39215-161f-4f21-b09f-3536e270730a/hudi-utilities-slim-bundle_2.12-0.14.0.jar
 to class loader
   23/11/26 09:24:54 INFO Utils: Successfully started service 
'org.apache.spark.network.netty.NettyBlockTransferService' on port 64397.
   23/11/26 09:24:54 INFO NettyBlockTransferService: Server created on 
soumils-mbp:64397
   23/11/26 09:24:54 INFO BlockManager: Using 
org.apache.spark.storage.RandomBlockReplicationPolicy for block replication 
policy
   23/11/26 09:24:54 INFO BlockManagerMaster: Registering BlockManager 
BlockManagerId(driver, soumils-mbp, 64397, None)
   23/11/26 09:24:54 INFO BlockManagerMasterEndpoint: Registering block manager 
soumils-mbp:64397 with 434.4 MiB RAM, BlockManagerId(driver, soumils-mbp, 
64397, None)
   23/11/26 09:24:54 INFO BlockManagerMaster: Registered BlockManager 
BlockManagerId(driver, soumils-mbp, 64397, None)
   23/11/26 09:24:54 INFO BlockManager: Initialized BlockManager: 
BlockManagerId(driver, soumils-mbp, 64397, None)
   23/11/26 09:24:54 WARN DFSPropertiesConfiguration: Cannot find 
HUDI_CONF_DIR, please set it as the dir of hudi-defaults.conf
   23/11/26 09:24:54 WARN DFSPropertiesConfiguration: Properties file 
file:/etc/hudi/conf/hudi-defaults.conf not found. Ignoring to load props file
   23/11/26 09:24:54 INFO SharedState: Setting hive.metastore.warehouse.dir 
('null') to the value of spark.sql.warehouse.dir.
   23/11/26 09:24:54 INFO SharedState: Warehouse path is 
'file:/Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/spark-warehouse'.
   23/11/26 09:24:55 INFO HoodieStreamer: Creating Hudi Streamer with configs:
   auto.offset.reset: earliest
   bootstrap.servers: localhost:7092
   hoodie.auto.adjust.lock.configs: true
   hoodie.datasource.write.partitionpath.field: order_date
   hoodie.datasource.write.precombine.field: ts
   hoodie.datasource.write.reconcile.schema: false
   hoodie.datasource.write.recordkey.field: order_id
   hoodie.deltastreamer.schemaprovider.registry.url: 
http://localhost:8081/subjects/postgres.public.sales-value/versions/latest
   hoodie.deltastreamer.source.kafka.topic: postgres.public.sales
   hoodie.deltastreamer.source.kafka.value.deserializer.class: 
org.apache.hudi.utilities.deser.KafkaAvroSchemaDeserializer
   schema.registry.url: http://localhost:8081
   
   23/11/26 09:24:55 INFO HoodieSparkKeyGeneratorFactory: The value of 
hoodie.datasource.write.keygenerator.type is empty; inferred to be SIMPLE
   23/11/26 09:24:55 INFO HoodieSparkKeyGeneratorFactory: The value of 
hoodie.datasource.write.keygenerator.type is empty; inferred to be SIMPLE
   23/11/26 09:24:55 INFO HoodieTableMetaClient: Initializing 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
 as hoodie table 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
   23/11/26 09:24:55 INFO HoodieTableMetaClient: Loading HoodieTableMetaClient 
from 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
   23/11/26 09:24:55 INFO HoodieTableConfig: Loading table properties from 
file:/Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders/.hoodie/hoodie.properties
   23/11/26 09:24:55 INFO HoodieTableMetaClient: Finished Loading Table of type 
COPY_ON_WRITE(version=1, baseFileFormat=PARQUET) from 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
   23/11/26 09:24:55 INFO HoodieTableMetaClient: Finished initializing Table of 
type COPY_ON_WRITE from 
file:///Users/soumilshah/IdeaProjects/SparkProject/apache-hudi-delta-streamer-labs/E7/hudidb/orders
   23/11/26 09:24:55 WARN ConfigUtils: The configuration key 
'hoodie.deltastreamer.schemaprovider.registry.url' has been deprecated and may 
be removed in the future. Please use the new key 
'hoodie.streamer.schemaprovider.registry.url' instead.
   23/11/26 09:24:55 INFO SparkContext: SparkContext is stopping with exitCode 
0.
   23/11/26 09:24:55 INFO SparkUI: Stopped Spark web UI at 
http://soumils-mbp:8090
   23/11/26 09:24:55 INFO MapOutputTrackerMasterEndpoint: 
MapOutputTrackerMasterEndpoint stopped!
   23/11/26 09:24:55 INFO MemoryStore: MemoryStore cleared
   23/11/26 09:24:55 INFO BlockManager: BlockManager stopped
   23/11/26 09:24:55 INFO BlockManagerMaster: BlockManagerMaster stopped
   23/11/26 09:24:55 INFO 
OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: 
OutputCommitCoordinator stopped!
   23/11/26 09:24:55 INFO SparkContext: Successfully stopped SparkContext
   Exception in thread "main" 
org.apache.hudi.utilities.exception.HoodieSchemaFetchException: Error reading 
source schema from registry 
:http://localhost:8081/subjects/postgres.public.sales-value/versions/latest
           at 
org.apache.hudi.utilities.schema.SchemaRegistryProvider.getSourceSchema(SchemaRegistryProvider.java:198)
           at 
org.apache.hudi.utilities.streamer.StreamSync.registerAvroSchemas(StreamSync.java:1104)
           at 
org.apache.hudi.utilities.streamer.StreamSync.<init>(StreamSync.java:290)
           at 
org.apache.hudi.utilities.streamer.HoodieStreamer$StreamSyncService.<init>(HoodieStreamer.java:707)
           at 
org.apache.hudi.utilities.streamer.HoodieStreamer.<init>(HoodieStreamer.java:159)
           at 
org.apache.hudi.utilities.streamer.HoodieStreamer.<init>(HoodieStreamer.java:131)
           at 
org.apache.hudi.utilities.streamer.HoodieStreamer.main(HoodieStreamer.java:584)
           at 
java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
           at 
java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
           at 
java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
           at java.base/java.lang.reflect.Method.invoke(Method.java:566)
           at 
org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52)
           at 
org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:1020)
           at 
org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:192)
           at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:215)
           at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:91)
           at 
org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1111)
           at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1120)
           at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
   Caused by: org.apache.hudi.internal.schema.HoodieSchemaException: Failed to 
parse schema from registry: 
{"type":"record","name":"Envelope","namespace":"postgres.public.sales","fields":[{"name":"before","type":["null",{"type":"record","name":"Value","fields":[{"name":"salesid","type":"int"},{"name":"invoiceid","type":["null","int"],"default":null},{"name":"itemid","type":["null","int"],"default":null},{"name":"category","type":["null","string"],"default":null},{"name":"price","type":["null",{"type":"bytes","scale":2,"precision":10,"connect.version":1,"connect.parameters":{"scale":"2","connect.decimal.precision":"10"},"connect.name":"org.apache.kafka.connect.data.Decimal","logicalType":"decimal"}],"default":null},{"name":"quantity","type":["null","int"],"default":null},{"name":"orderdate","type":["null",{"type":"int","connect.version":1,"connect.name":"io.debezium.time.Date"}],"default":null},{"name":"destinationstate","type":["null","string"],"default":null},{"name":"shippingtype",
 
"type":["null","string"],"default":null},{"name":"referral","type":["null","string"],"default":null},{"name":"updated_at","type":["null",{"type":"long","connect.version":1,"connect.name":"io.debezium.time.MicroTimestamp"}],"default":null}],"connect.name":"postgres.public.sales.Value"}],"default":null},{"name":"after","type":["null","Value"],"default":null},{"name":"source","type":{"type":"record","name":"Source","namespace":"io.debezium.connector.postgresql","fields":[{"name":"version","type":"string"},{"name":"connector","type":"string"},{"name":"name","type":"string"},{"name":"ts_ms","type":"long"},{"name":"snapshot","type":[{"type":"string","connect.version":1,"connect.parameters":{"allowed":"true,last,false"},"connect.default":"false","connect.name":"io.debezium.data.Enum"},"null"],"default":"false"},{"name":"db","type":"string"},{"name":"schema","type":"string"},{"name":"table","type":"string"},{"name":"txId","type":["null","long"],"default":null},{"name":"lsn","type":["null","
 
long"],"default":null},{"name":"xmin","type":["null","long"],"default":null}],"connect.name":"io.debezium.connector.postgresql.Source"}},{"name":"op","type":"string"},{"name":"ts_ms","type":["null","long"],"default":null}],"connect.name":"postgres.public.sales.Envelope"}
           at 
org.apache.hudi.utilities.schema.SchemaRegistryProvider.parseSchemaFromRegistry(SchemaRegistryProvider.java:105)
           at 
org.apache.hudi.utilities.schema.SchemaRegistryProvider.getSourceSchema(SchemaRegistryProvider.java:196)
           ... 18 more
   Caused by: java.lang.IllegalArgumentException: Property 
hoodie.streamer.schemaprovider.registry.schemaconverter not found
           at 
org.apache.hudi.common.util.ConfigUtils.getStringWithAltKeys(ConfigUtils.java:334)
           at 
org.apache.hudi.common.util.ConfigUtils.getStringWithAltKeys(ConfigUtils.java:308)
           at 
org.apache.hudi.utilities.schema.SchemaRegistryProvider.parseSchemaFromRegistry(SchemaRegistryProvider.java:99)
           ... 19 more
   23/11/26 09:24:55 INFO ShutdownHookManager: Shutdown hook called
   23/11/26 09:24:55 INFO ShutdownHookManager: Deleting directory 
/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-b5b317b5-a47b-4ed9-8f22-db44bc308194
   23/11/26 09:24:55 INFO ShutdownHookManager: Deleting directory 
/private/var/folders/qq/s_1bjv516pn_mck29cwdwxnm0000gp/T/spark-3f9f543b-ad9d-459b-95cc-19fc7f2573e0
   (venv) soumilshah@Soumils-MBP E7 % 
   
   ```
   If anyone in the community has experience with Delta Streamer and Debezium, 
and can offer insights or guidance on potential misconfigurations, I would 
greatly appreciate your assistance. I'm eager to learn and ensure the success 
of my Delta Streamer job.
   
   Thank you in advance for your time and expertise.
   
   Best regards,
   
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to