This message was deleted.
# general
s
This message was deleted.
y
this is on postgres side, can you share your pg table schema?
d
Agreed. The error indicates that Postgres is throwing an exception. However, the UUID field doesn’t appear to be getting properly parsed inside the connector.
f
@Yu Wei Sung my postgres side is like below. so,
uuid
is actually a primary key. since it's part of the requirment.
Copy code
CREATE TABLE public.jdbc_sink (
	uuid uuid NOT NULL,
	"timestamp" int8 NULL,
	CONSTRAINT jdbc_sink_pkey PRIMARY KEY (uuid)
);
@David K yes, actually i'm not sure my json (
UUID
and
timestamp
) is being read by the java file. is there anyway i can reproduce how the sink work? i found the code in the github. but i cannot see the process one by one. e.g. debug the parsed value on each line. https://github.com/apache/pulsar/blob/master/pulsar-io/jdbc/core/src/main/java/org/apache/pulsar/io/jdbc/JdbcAbstractSink.java
y
you can try turn on debug mode and trace this line. https://github.com/apache/pulsar/blob/a9d5d25e52710503cc2642ba7a8aadda0a32faae/pul[…]main/java/org/apache/pulsar/io/jdbc/BaseJdbcAutoSchemaSink.java my suspicion is the data type. in your avro schem, uuid: string but in pg table, it is uuid: uuid. easy troubleshooting step: you can try another table with uuid: text as primary key.
👍 1
Copy code
{
  "type": "AVRO",
  "schema": "{\"type\":\"record\",\"name\":\"jdbc_sink\",\"fields\":[{\"name\":\"uuid\",\"type\": {\"type\": \"string\", \"logicalType\": \"uuid\"},{\"name\":\"timestamp\",\"type\":\"long\"}]}",
  "properties": {}
}
try this.
f
Thanks @Yu Wei Sung, it was thoughtful one that the error was on schema. however i already try it, but still getting the same error. The schema using
pulsar-admin schema get
function.
Copy code
root@d5da1b9bf4c2:/pulsar# bin/pulsar-admin schemas get pulsar-postgres-jdbc-sink-topic      
{
  "version": 2,
  "schemaInfo": {
    "name": "jdbc_source_output",
    "schema": {
      "type": "record",
      "name": "jdbc_sink",
      "fields": [
        {
          "name": "uuid",
          "type": {
            "type": "string",
            "logicalType": "uuid"
          }
        },
        {
          "name": "timestamp",
          "type": "long"
        }
      ]
    },
    "type": "AVRO",
    "timestamp": 1697675942684,
    "properties": {}
  }
}
my error logs:
Copy code
2023-10-19T07:01:31,777+0000 [pulsar-timer-9-1] INFO  org.apache.pulsar.client.impl.NegativeAcksTracker - [ConsumerBase{subscription='public/default/pulsar-postgres-jdbc-sink', consumerName='e07e2', topic='jdbc_source_output'}] 5 messages will be re-delivered
2023-10-19T07:01:32,248+0000 [pool-5-thread-1] ERROR org.apache.pulsar.io.jdbc.JdbcAbstractSink - Got exception ERROR: null value in column "uuid" of relation "jdbc_sink" violates not-null constraint
  Detail: Failing row contains (null, null). after 2 ms, failing 5 messages
org.postgresql.util.PSQLException: ERROR: null value in column "uuid" of relation "jdbc_sink" violates not-null constraint
  Detail: Failing row contains (null, null).
	at org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2676) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2366) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:356) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:496) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.jdbc.PgStatement.execute(PgStatement.java:413) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.jdbc.PgPreparedStatement.executeWithFlags(PgPreparedStatement.java:190) ~[postgresql-42.5.1.jar:42.5.1]
	at org.postgresql.jdbc.PgPreparedStatement.execute(PgPreparedStatement.java:177) ~[postgresql-42.5.1.jar:42.5.1]
	at org.apache.pulsar.io.jdbc.JdbcAbstractSink.flush(JdbcAbstractSink.java:289) ~[pulsar-io-jdbc-core-3.1.0.jar:?]
	at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) ~[?:?]
	at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:305) ~[?:?]
	at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[?:?]
	at java.lang.Thread.run(Thread.java:833) ~[?:?]
2023-10-19T07:02:31,611+0000 [pulsar-timer-9-1] INFO  org.apache.pulsar.client.impl.ConsumerStatsRecorderImpl - [jdbc_source_output] [public/default/pulsar-postgres-jdbc-sink] [e07e2] Prefetched messages: 0 --- Consume throughput received: 0.08 msgs/s --- 0.00 Mbit/s --- Ack sent rate: 0.00 ack/s --- Failed messages: 0 --- batch messages: 0 ---Failed acks: 0
2023-10-19T07:02:51,809+0000 [pulsar-timer-9-1] INFO  org.apache.pulsar.client.impl.NegativeAcksTracker - [ConsumerBase{subscription='public/default/pulsar-postgres-jdbc-sink', consumerName='e07e2', topic='jdbc_source_output'}] 5 messages will be re-delivered
will try the debug mode for the time being.
y
report this to gh issue if this doesn’t work.
f
@Yu Wei Sung @David K i use some altrenative by create some
pulsar function
with python. in summary this pulsar function use
psycopg2
to connect directly to the db and parse some json from the topics, and use
INSERT INTO {table_name} ({columns}) VALUES ({placeholders})
CRUD to insert.
Copy code
import pulsar
from pulsar import Function
import psycopg2
import json

class PostgresSync(Function):

    def process(self, input, context):
        error_messages = []

        context.get_logger().info(f"Received input: {input}")
        
        conn = None
        cursor = None
        
        try:
            context.get_logger().info("Trying to establish PostgreSQL connection...")
            
            conn = psycopg2.connect(
                host="localhost",
                port=5432,
                dbname="postgres",
                user="postgres",
                password="password"
            )
            cursor = conn.cursor()
            
            context.get_logger().info("Successfully established PostgreSQL connection.")
        except Exception as e:
            err_msg = f"Error while establishing connection to PostgreSQL: {e}"
            context.get_logger().error(err_msg)
            error_messages.append(err_msg)

        try:
            # Parse the content
            data = json.loads(input)
            after_data = data.get("after", {})
            
            # Extract table name and append "_destination"
            table_name = data.get("source", {}).get("table", "") + "_destination"

            # Constructing dynamic query
            columns = ", ".join(after_data.keys())
            placeholders = ", ".join(["%s"] * len(after_data))
            values = tuple(after_data.values())
            
            query = f"INSERT INTO {table_name} ({columns}) VALUES ({placeholders});"
            
            cursor.execute(query, values)
            conn.commit()
            
            context.get_logger().info("Successfully inserted data into PostgreSQL.")
        except Exception as e:
            err_msg = f"Error during data insertion into PostgreSQL. Data: {input}. Error: {e}"
            context.get_logger().error(err_msg)
            error_messages.append(err_msg)

        try:
            # Close the connection
            cursor.close()
            conn.close()
            context.get_logger().info("Closed PostgreSQL connection successfully.")
        except Exception as e:
            err_msg = f"Error while closing the PostgreSQL connection: {e}"
            context.get_logger().error(err_msg)
            error_messages.append(err_msg)

        return "; ".join(error_messages) if error_messages else input
however i do need some advise if this code will scalable for thousands to hundred thousands row inputs or not.
d
This approach is not scalable. For starters recreating the database connection every record is inefficient. The connection should be created once and then cached. Secondly the record should be published in a batch rather than individually. However even these changes will not work for 100k+ rows per second.
👍 1
y
Architecture wise, it is jdbc, aka connection. You should not use multiple jdbc connections to sink to a single table. Both pulsar and postgres have their own transaction mechanisms. There are many scalability concerns coordinating those rollback commits acknowledge….to consider
👍 1
f
sure, thanks for the input. already discard this option. right now i try to understand better the pulsar schema. both AVRO and JSON.