Slackbot
10/16/2023, 4:24 AMYu Wei Sung
10/16/2023, 2:50 PMDavid K
10/16/2023, 3:28 PMFirman Insan Muhammad
10/18/2023, 4:27 AMuuid is actually a primary key. since it's part of the requirment.
CREATE TABLE public.jdbc_sink (
uuid uuid NOT NULL,
"timestamp" int8 NULL,
CONSTRAINT jdbc_sink_pkey PRIMARY KEY (uuid)
);Firman Insan Muhammad
10/18/2023, 4:31 AMUUID 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.javaYu Wei Sung
10/18/2023, 12:41 PMYu Wei Sung
10/18/2023, 12:52 PM{
"type": "AVRO",
"schema": "{\"type\":\"record\",\"name\":\"jdbc_sink\",\"fields\":[{\"name\":\"uuid\",\"type\": {\"type\": \"string\", \"logicalType\": \"uuid\"},{\"name\":\"timestamp\",\"type\":\"long\"}]}",
"properties": {}
}
try this.Firman Insan Muhammad
10/19/2023, 7:04 AMpulsar-admin schema get function.
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": {}
}
}Firman Insan Muhammad
10/19/2023, 7:06 AM2023-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-deliveredFirman Insan Muhammad
10/19/2023, 7:18 AMYu Wei Sung
10/19/2023, 12:56 PMFirman Insan Muhammad
10/22/2023, 9:04 AMpulsar 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.
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 inputFirman Insan Muhammad
10/22/2023, 9:05 AMDavid K
10/23/2023, 3:24 PMYu Wei Sung
10/23/2023, 4:04 PMFirman Insan Muhammad
11/01/2023, 5:33 PM