Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion clients/python/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ version = "0.20.12"
description = "Taskbroker python client and worker runtime"
readme = "README.md"
dependencies = [
"sentry-arroyo>=2.41.0",
"sentry-arroyo>=2.41.1",
"sentry-sdk[http2]>=2.43.0",
"sentry-protos>=0.26.1",
"confluent_kafka>=2.3.0",
Expand Down
5 changes: 2 additions & 3 deletions clients/python/src/examples/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,13 @@
from time import sleep
from typing import Any

from arroyo.backends.kafka import KafkaPayload, KafkaProducer
from arroyo.backends.kafka import FutureTrackingProducer, KafkaPayload, KafkaProducer
from arroyo.types import Topic
from redis import StrictRedis

from examples.app import app
from taskbroker_client.retry import LastAction, NoRetriesRemainingError, Retry, RetryTaskError
from taskbroker_client.retry import retry_task as retry_task_helper
from taskbroker_client.worker.producer import TaskProducer
from taskbroker_client.worker.workerchild import ProcessingDeadlineExceeded

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -135,7 +134,7 @@ def task_that_produces(
def producer_factory() -> KafkaProducer:
return KafkaProducer({"bootstrap.servers": bootstrap_servers})

producer = TaskProducer("test.producer", producer_factory)
producer = FutureTrackingProducer("test.producer", producer_factory)
production_count = random.randint(1, 50) if random_count else production_count
for i in range(production_count):
logger.debug(f"Producing message {i} onto topic {destination_topic}...")
Expand Down
105 changes: 0 additions & 105 deletions clients/python/src/taskbroker_client/worker/producer.py

This file was deleted.

8 changes: 2 additions & 6 deletions clients/python/src/taskbroker_client/worker/workerchild.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,7 @@
import sentry_sdk
import zstandard as zstd
from arroyo.backends.abstract import ProducerFuture
from arroyo.backends.kafka import KafkaPayload
from arroyo.backends.kafka.producer import FutureTrackingProducer
from arroyo.backends.kafka import FutureTrackingProducer, KafkaPayload
from arroyo.types import BrokerValue
from sentry_protos.taskbroker.v1.taskbroker_pb2 import (
TASK_ACTIVATION_STATUS_COMPLETE,
Expand All @@ -39,7 +38,6 @@
from taskbroker_client.state import clear_current_task, current_task, set_current_task
from taskbroker_client.task import Task
from taskbroker_client.types import ContextHook, InflightTaskActivation, ProcessingResult
from taskbroker_client.worker.producer import TaskProducer

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -547,9 +545,7 @@ def check_task_future_completion(

# To have Taskworker track futures, set the env var `ARROYO_TRACK_PRODUCER_FUTURES = True`
# in the worker process
task_produced_futures = (
TaskProducer.collect_futures() | FutureTrackingProducer.collect_futures()
)
task_produced_futures = FutureTrackingProducer.collect_futures()

# If the task function itself failed, we don't need to await any
# producer futures since it'll be retried anyways
Expand Down
96 changes: 0 additions & 96 deletions clients/python/tests/worker/test_producer.py

This file was deleted.

Loading
Loading