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]
