apache / apache/pulsar-client-python

[BUG] use celery task send data to pulsar, init pulsar client error

未關閉
#150 3 則留言 0 個 reaction 已指派 0 人 在 GitHub 檢視
主要語言
Python
星號
75
分支
53
PR 合併指標
30 天內沒有已合併 PR

描述

python: 3.6
pulsar-client: pulsar-client[avro]==2.10.2
celery: 5.1.2

code example:
```python
#!/usr/bin/env python

import time
import random
import string
from pulsar import Client, CompressionType
from pulsar.schema import AvroSchema, Record, String, Integer

def generate_random_string(length=6):
charset = string.ascii_letters + string.digits
random_chars = random.choices(charset, k=length)
random_string = "".join(random_chars)
return random_string.capitalize()

class User(Record):
name = String()
age = Integer

UserAvroSchema = AvroSchema(User) # type: ignore

def gen_random_data():
return User(user=generate_random_string(), age=random.randint(0, 100))

class PulsarDemo(object):
def __init__(self) -> None:
self.SERVICE_URL = "pulsar://***"
self.TOPIC = "persistent://****"
client = Client(service_url=self.SERVICE_URL)
self.producer = client.create_producer(
topic=self.TOPIC,
schema=UserAvroSchema,
batching_enabled=True,
batching_max_messages=1000,
batching_max_publish_delay_ms=1000,
compression_type=CompressionType.SNAPPY, # type: ignore
)

def send_callback(self, send_result, msg_id):
print("Message published: result:{} msg_id:{}".format(send_result, msg_id))

def async_producer(self, cnt=1000):
while cnt >= 0:
data = gen_random_data()
self.producer.send_async(
data,
callback=self.send_callback,
)
time.sleep(0.01)
cnt -= 1
self.producer.flush()

# celery task
from celery import shared_task
@shared_task
def mock_data2pulsar(cnt=1000):
mock = PulsarDemo()
mock.async_producer()

```

an exception occurred at:
```text
[2023-08-29 10:27:46,933: ERROR/ForkPoolWorker-31] Pulsar error: TopicNotFound
```

貢獻指南

開啟貢獻指南

研究方向

從提供的 Celery task 和 PulsarDemo.__init__ 路徑開始,特別檢查 Client 和 create_producer,並使用列出的 Python、pulsar-client 和 Celery 版本重現回報的 TopicNotFound 錯誤。完成的標準是確認該故障是由 task/client 初始化還是 topic 設定所造成,並記錄預期行為與經驗證的解決方案。

由索引模型根據 Issue 內容生成。

評估

技術堆疊
python
領域
distributed-systems
Issue 類型
缺陷
難度
4/5
預估耗時
3-5 天
活躍度
停滯
描述清晰度
需要釐清
新手友好度
25/100

把新 issue 寄到你的電子郵件信箱

精選適合新手參與的 GitHub issue 摘要。