Slackbot
10/04/2023, 9:02 PMDavid K
10/04/2023, 9:11 PMschemas get command AFTER you tried to manually create the topic?Oneeb
10/04/2023, 9:36 PMOneeb
10/04/2023, 9:37 PM{
"version": 2,
"schemaInfo": {
"name": "customer-processed",
"schema": {
"type": "record",
"name": "User",
"namespace": "data-src-conn",
"fields": [
{
"name": "id",
"type": [
"null",
"string"
]
},
{
"name": "status",
"type": [
"null",
"string"
]
},
{
"name": "status_metadata",
"type": [
"string",
"null"
]
},
{
"name": "create_metadata",
"type": [
"string",
"null"
]
},
{
"name": "update_metadata",
"type": [
"string",
"null"
]
}
]
},
"type": "JSON",
"timestamp": 1696426118370,
"properties": {}
}
}
I have attached the definition config file I used:Oneeb
10/04/2023, 9:38 PMDavid K
10/04/2023, 10:14 PMOneeb
10/04/2023, 10:19 PMDavid K
10/04/2023, 10:20 PMcustomSchemaOutputs property in the configs.David K
10/04/2023, 10:25 PMoutputSchemaType property to “JSON”. See the REST API for a detailed example.Oneeb
10/04/2023, 10:35 PMOneeb
10/05/2023, 8:30 AMOneeb
10/05/2023, 8:40 AMOneeb
10/05/2023, 2:59 PMReceived error from server: Key schemas or Value schemas are different schema type, from key schema type is BYTES and to key schema is BYTES, from value schema is JSON and to value schema is BYTES
How can I resolve this?David K
10/05/2023, 4:16 PMDavid K
10/05/2023, 4:20 PMOneeb
10/05/2023, 4:21 PM{
"before": {
"id": "f7a5d1df-7a7b-4ebe-8262-ab5078fb056f",
"status": "INACTIVE",
"status_metadata": "\"\"",
"create_metadata": "{\"time\": 1696517856.219328, \"actor\": \"f7a5d1df-7a7b-4ebe-8262-ab5078fb056f\", \"actorType\": \"CUSTOMER\", \"callerService\": \"\"}",
"update_metadata": "{\"time\": 1696517856.219328, \"actor\": \"f7a5d1df-7a7b-4ebe-8262-ab5078fb056f\", \"actorType\": \"CUSTOMER\", \"callerService\": \"\"}"
},
"after": null,
"source": {
"version": "1.9.7.Final",
"connector": "postgresql",
"name": "dbserver101",
"ts_ms": 1696517858259,
"snapshot": "false",
"db": "customer",
"sequence": "[\"36248032\",\"36248032\"]",
"schema": "customer",
"table": "customers",
"txId": 1542,
"lsn": 36248032,
"xmin": null
},
"op": "d",
"ts_ms": 1696517858479,
"transaction": null
}Oneeb
10/05/2023, 4:22 PMlocalrunOneeb
10/05/2023, 4:23 PMoutputSchemaType and customSchemaOutputs . Thats why Im now writing a Java Pulsar Function, and Im using the ExclamationFunction example just to see how it works.David K
10/05/2023, 4:27 PMpulsar-admin schema get command for the topic?Oneeb
10/05/2023, 4:28 PMOneeb
10/05/2023, 4:29 PM{
"version": 1,
"schemaInfo": {
"name": "customer-processed",
"schema": "\"string\"",
"type": "AVRO",
"timestamp": 1696519395927,
"properties": {
"__alwaysAllowNull": "true",
"__jsr310ConversionEnabled": "false"
}
}
}David K
10/05/2023, 4:29 PMOneeb
10/05/2023, 4:30 PMpulsar-admin schema getDavid K
10/05/2023, 4:31 PMOneeb
10/05/2023, 4:32 PMDavid K
10/05/2023, 4:33 PMDavid K
10/05/2023, 4:34 PMschema get command for the first topic. I only see the second one, which is an AVRO schemaOneeb
10/05/2023, 4:34 PMDavid K
10/05/2023, 4:49 PMinputSpecs: - "input-topic-name": schemaType: JSON , as defined here.Oneeb
10/05/2023, 4:50 PMbin/pulsar-admin functions localrun \
--jar $PF_JAR_PATH \
--className org.example.ExclamationFunction \
--inputs <persistent://public/data-src-conn/dbserver101.customer.customers> \
--output <persistent://public/data-src-conn/customer-processed> \
--name ExclamationFunction \
--schema-type json
This is already how Im running the function ^Oneeb
10/05/2023, 5:10 PMbin/pulsar-admin functions localrun --jar $PF_JAR_PATH --function-config-file $PF_CONFIG_PATH and the yaml configs are as follows
className: "org.example.ExclamationFunction"
inputs:
- "<persistent://public/data-src-conn/dbserver101.customer.customers>"
output: "<persistent://public/data-src-conn/customer-processed>"
name: "ExclamationFunction"
schemaType: "json"
inputSpecs:
schemaType: "json"
And Im still getting the same error as beforeDavid K
10/05/2023, 5:13 PMDavid K
10/05/2023, 5:17 PMOneeb
10/05/2023, 5:18 PMclassName: "org.example.ExclamationFunction"
inputs:
- "<persistent://public/data-src-conn/dbserver101.customer.customers>"
output: "<persistent://public/data-src-conn/customer-processed>"
name: "ExclamationFunction"
schemaType: "json"
inputSpecs:
"<persistent://public/data-src-conn/dbserver101.customer.customers>":
schemaType: "json"
which now leads me to the error: Incompatible schema: exists schema type KEY_VALUE, new schema type JSONDavid K
10/05/2023, 5:21 PM"type": "KEY_VALUE", with the value being JSON.David K
10/05/2023, 5:22 PMDavid K
10/05/2023, 5:30 PMDavid K
10/05/2023, 5:30 PMSchema<KeyValue<Integer, String>> kvSchema = Schema.KeyValue(
Schema.BYTES,
Schema.JSON(YOUR_CLASS),
KeyValueEncodingType.SEPARATED
);
inputSpecs.put(IN, ConsumerConfig.builder().schemaType(
kvSchema.getSchemaInfo().getType().toString()).build());David K
10/05/2023, 5:31 PMOneeb
10/05/2023, 5:33 PMOneeb
10/05/2023, 5:35 PMbin/pulsar-admin localrun?David K
10/05/2023, 6:48 PMOneeb
10/05/2023, 6:49 PMOneeb
10/12/2023, 11:20 AMCustomerTableObject and returned that to Topic 3.
The schema for topic 3 now looks like:
{
"version": 0,
"schemaInfo": {
"name": "customer-processed-2",
"schema": {
"type": "record",
"name": "CustomerTableObject",
"namespace": "org.example",
"fields": [
{
"name": "create_metadata",
"type": [
"null",
"string"
]
},
{
"name": "id",
"type": [
"null",
"string"
]
},
{
"name": "status",
"type": [
"null",
"string"
]
},
{
"name": "status_metadata",
"type": [
"null",
"string"
]
},
{
"name": "update_metadata",
"type": [
"null",
"string"
]
}
]
},
"type": "JSON",
"timestamp": 1697034589493,
"properties": {
"__alwaysAllowNull": "true",
"__jsr310ConversionEnabled": "false"
}
}
}
When I run the Sink Connector now, I can see the _delta_log directory being created and a parquet file of size 0 bytes which contains the schema.
• When I query this delta table using Spark, I can only see the schema
However I don’t see any of the payload being written down. When I send messages to topic 3 with the following payload:
{
"id": "apple",
"status": "banana",
"status_metadata": "caterpillar",
"create_metadata": "dog",
"update_metadata": "elep"
}
Nothing gets added.
However, when I run bin/pulsar-admin sinks status --tenant public --namespace data-src-conn --name delta_sink
I get this output:
{
"numInstances": 1,
"numRunning": 1,
"instances": [
{
"instanceId": 0,
"status": {
"running": true,
"error": "",
"numRestarts": 0,
"numReadFromPulsar": 6,
"numSystemExceptions": 0,
"latestSystemExceptions": [],
"numSinkExceptions": 0,
"latestSinkExceptions": [],
"numWrittenToSink": 6,
"lastReceivedTime": 1697106408287,
"workerId": "c-standalone-fw-localhost-8080"
}
}
]
}
I would expect numReadFromPulsar and/or numWrittenToSink to be 0 if its not successfully writing to my local file destination.
Just for greater clarity, if I query the topic status using bin/pulsar-admin topics stats <persistent://public/data-src-conn/topic-3> I get the following:
"public/data-src-conn/delta_sink" : {
"msgRateOut" : 0.0,
"msgThroughputOut" : 0.0,
"bytesOutCounter" : 972,
"msgOutCounter" : 6,
"msgRateRedeliver" : 0.0,
"messageAckRate" : 0.0,
"chunkedMessageRate" : 0,
"msgBacklog" : 3,
"backlogSize" : 1035,
"earliestMsgPublishTimeInBacklog" : 0,
"msgBacklogNoDelayed" : 3,
"blockedSubscriptionOnUnackedMsgs" : false,
"msgDelayed" : 0,
"unackedMessages" : 0,
"type" : "Failover",
"activeConsumerName" : "7d6e7",
"msgRateExpired" : 0.0,
"totalMsgExpired" : 0,
"lastExpireTimestamp" : 0,
"lastConsumedFlowTimestamp" : 1697106386342,
"lastConsumedTimestamp" : 1697106408286,
"lastAckedTimestamp" : 0,
"lastMarkDeleteAdvancedTimestamp" : 0,
"consumers" : [ {
"msgRateOut" : 0.0,
"msgThroughputOut" : 0.0,
"bytesOutCounter" : 972,
"msgOutCounter" : 6,
"msgRateRedeliver" : 0.0,
"messageAckRate" : 0.0,
"chunkedMessageRate" : 0.0,
"consumerName" : "7d6e7",
"availablePermits" : 994,
"unackedMessages" : 0,
"avgMessagesPerEntry" : 2,
"blockedConsumerOnUnackedMsgs" : false,
"address" : "/127.0.0.1:57742",
"connectedSince" : "2023-10-12T13:26:26.339217+03:00",
"clientVersion" : "Pulsar-Java-v3.1.0",
"lastAckedTimestamp" : 0,
"lastConsumedTimestamp" : 1697106408286,
"lastConsumedFlowTimestamp" : 1697106386342,
"metadata" : {
"instance_id" : "0",
"application" : "pulsar-sink",
"instance_hostname" : "Oneebs-MacBook-Pro.local",
"id" : "public/data-src-conn/delta_sink"
},
"lastAckedTime" : "1970-01-01T04:00:00+04:00",
"lastConsumedTime" : "2023-10-12T13:26:48.286+03:00"
}
What am I missing? Why isn’t the Sink Connector writing delta parquet files to my local path?Oneeb
10/12/2023, 11:24 AMDavid K
10/12/2023, 5:02 PMOneeb
10/12/2023, 8:58 PMOneeb
10/13/2023, 9:02 AMOneeb
10/17/2023, 4:54 AMHang Chen
10/20/2023, 3:59 AM