This message was deleted.
# general
s
This message was deleted.
r
what can be the issue in this config ?
@Alexander Preuss could you please provide some help ?
a
@Rajnesh I haven’t worked with the Delta Lake connector before. Maybe @Hang Chen can help
h
Oh sorry, I will take a look
How many messages have you produced to the topic?
I’m not sure if it is caused by sinker not commit to the delta lake
r
I have tried with one message and also tried with multiple messages on one topic
here is the code through which I am pushing the data
Copy code
from pulsar.schema import *
import pulsar
import time
from dataclasses import dataclass
from typing import List
import uuid
from datetime import datetime


def post_data():
    class AuditEvent(Record):
        event = String()
        description = String()
        affected_id= String()
        affected_type= String()
        actor_id=String()
        actor_type=String()
    
    print("schema creation in process")
    client = pulsar.Client('<pulsar://localhost:6650>')
    time.sleep(5)

    producer = client.create_producer(
        topic='public/default/audit_file',
        schema=JsonSchema(AuditEvent),
    )

    print("post_data")
    producer.send(AuditEvent(
                            event='audit', 
                            description='Creating Event',
                            affected_id='328006c1-793f-43be-93824906',
                            affected_type='ACTOR_TYPE_SERVICE',
                            actor_id='018ac9c5-f594-722f-aeda7524b7a',
                            actor_type='ACTOR_TYPE_SERVICE'
                            )
                )
    
    
    print("data posted")


post_data()
if I run it again and again, it produces a new empty parquet file for each run. this is the screenshot of the directory,
h
Ok, I’m trying to reproduce it
r
sure
@Hang Chen , did you reproduce the error ?
@Hang Chen, any recommendations ?
h
Sorry for the late response
I will provide the details this night or tomorrow
👍 1
@Yan Zhao
y
Ok, I will trace this problem
👍 1
r
@Yan Zhao
y
Sry, I’m trace it now.
👍 1
Could you wait for
120s
, the delta writer commit interval is
120s
. After commit, the parquet file will be flushed to file system.
r
Could you please elaborate more, should i add delay operation in my code ?
y
Ok, the sink writer received the msg, and then write it to the parquet file. After commit, then the parquet file data will be flushed to the filesystem. When does the parquet file commit? After the sink writer write the msg to the parquet file, it will check two condition. 1. Whether the msg number of current parquet file is more than
maxRecordsPerCommit
, if so, commit. 2. Whether the time between the current time and the last commit exceeds
maxCommitInterval
, if so, commit. So I guess that the parquet file haven’t commit yet, you can start the sink writer, and continuously sending messages to the topic for more than 120 seconds. Then check the parquet file status. If possible, could you attach the logs when you test it. Thanks.
h
We can try to decrease
maxCommitInterval
to 10
r
sure Thanks Hang, I will try decreasing
maxCommitInterval
to 10
y
Hi, @Rajnesh. @Oneeb encounter the same issue with you. It’s a compatible problem with mac m1 core, @Oneeb will fix it.
👍 2
🙌 2