This is the multi-page printable view of this section. Click here to print.

Return to the regular view of this page.

Dapr Python SDK

用于开发 Dapr 应用程序的 Python SDK 软件包

Dapr 提供了多种子包来帮助开发 Python 应用程序。使用它们,你可以使用 Dapr 创建 Python 客户端、服务器和虚拟 actor。

前提条件

安装

要开始使用 Python SDK,请安装主要的 Dapr Python SDK 软件包。

pip install dapr

注意: 开发包将包含与 Dapr 运行时预发布版本兼容的功能和行为。在安装 dapr-dev 软件包之前,请确保卸载任何稳定版本的 Python SDK。

pip install dapr-dev

可用的子包

SDK 导入

Python SDK 导入是随主 SDK 安装包含的子包,但在使用时需要导入。Dapr Python SDK 提供的最常见导入包括:

Client

编写 Python 应用程序以与 Dapr 边车和其他 Dapr 应用程序交互,包括 Python 中的有状态虚拟 actor

Actors

创建并与 Dapr 的 Actor 框架交互。

Conversation

使用 Dapr Conversation API(Alpha)进行 LLM 交互、工具和多轮流程。

了解有关 所有 可用的 Dapr Python SDK 导入 的更多信息。

SDK 扩展

SDK 扩展主要用作接收发布订阅事件、以编程方式创建发布订阅订阅以及处理输入绑定事件的实用工具。虽然你可以在不使用扩展的情况下完成所有这些任务,但使用 Python SDK 扩展会更加方便。

gRPC

使用 gRPC 服务器扩展创建 Dapr 服务。

FastAPI

使用 Dapr FastAPI 扩展与 Dapr Python 虚拟 actor 和发布订阅集成。

Flask

使用 Dapr Flask 扩展与 Dapr Python 虚拟 actor 集成。

Workflow

在 Python 中编写与其他 Dapr API 配合使用的工作流。

了解有关 Dapr Python SDK 扩展 的更多信息。

试用

克隆 Python SDK 仓库。

git clone https://github.com/dapr/python-sdk.git

浏览 Python 快速入门、教程和示例,查看 Dapr 的实际运行情况:

SDK 示例描述
快速入门在几分钟内使用 Python SDK 体验 Dapr 的 API 构建块。
SDK 示例克隆 SDK 仓库以试用一些示例并开始使用。
绑定教程了解 Dapr Python SDK 如何与其他 Dapr SDK 配合工作以启用绑定。
分布式计算器教程使用 Dapr Python SDK 处理方法调用和状态持久化功能。
Hello World 教程了解如何使用 Python SDK 在本地机器上启动和运行 Dapr。
Hello Kubernetes 教程在 Kubernetes 集群中使用 Dapr Python SDK 启动和运行。
可观测性教程使用 Python SDK 探索 Dapr 的指标收集、链路追踪、日志记录和健康检查功能。
发布订阅教程了解 Dapr Python SDK 如何与其他 Dapr SDK 配合工作以启用发布订阅应用程序。

更多信息

序列化

了解有关 Dapr SDK 中序列化的更多信息。

PyPI

Python 软件包索引

1 - Dapr 客户端 Python SDK 入门

如何快速上手使用 Dapr Python SDK

Dapr 客户端包允许你从 Python 应用程序与其他 Dapr 应用程序进行交互。

前提条件

在开始之前,安装 Dapr Python 包

导入客户端包

dapr 包包含 DaprClient,用于创建和使用客户端。

from dapr.clients import DaprClient

初始化客户端

你可以通过多种方式初始化 Dapr 客户端:

默认值:

当不使用任何参数初始化客户端时,它将使用 Dapr 边车实例的默认值(127.0.0.1:50001)。

from dapr.clients import DaprClient

with DaprClient() as d:
    # 使用客户端

初始化时指定端点:

当在构造函数中作为参数传递时,gRPC 端点优先于任何配置或环境变量。

from dapr.clients import DaprClient

with DaprClient("mydomain:50051?tls=true") as d:
    # 使用客户端

配置选项:

Dapr 边车端点

你可以使用标准化的 DAPR_GRPC_ENDPOINT 环境变量来指定 gRPC 端点。设置此变量后,可以在不使用任何参数的情况下初始化客户端:

export DAPR_GRPC_ENDPOINT="mydomain:50051?tls=true"
from dapr.clients import DaprClient

with DaprClient() as d:
    # 客户端将使用环境变量中指定的端点

旧的环境变量 DAPR_RUNTIME_HOSTDAPR_HTTP_PORTDAPR_GRPC_PORT 也受支持,但 DAPR_GRPC_ENDPOINT 优先。

Dapr API 令牌

如果你的 Dapr 实例配置为需要 DAPR_API_TOKEN 环境变量,你可以在环境中设置它,客户端将自动使用它。
你可以在这里阅读更多关于 Dapr API 令牌身份验证的信息。

健康检查超时

在客户端初始化时,会针对 Dapr 边车(/healthz/outbound)执行健康检查。 客户端将等待边车启动并运行后再继续。

默认健康检查超时为 60 秒,但可以通过设置 DAPR_HEALTH_TIMEOUT 环境变量来覆盖。

重试和超时

如果从边车收到特定的错误代码,Dapr 客户端可以重试请求。这可以通过 DAPR_API_MAX_RETRIES 环境变量配置,并且会自动拾取,不需要任何代码更改。 DAPR_API_MAX_RETRIES 的默认值是 0,表示不会进行重试。

你可以通过创建 dapr.clients.retry.RetryPolicy 对象并将其传递给 DaprClient 构造函数来微调更多重试参数:

from dapr.clients.retry import RetryPolicy

retry = RetryPolicy(
    max_attempts=5, 
    initial_backoff=1, 
    max_backoff=20, 
    backoff_multiplier=1.5,
    retryable_http_status_codes=[408, 429, 500, 502, 503, 504],
    retryable_grpc_status_codes=[StatusCode.UNAVAILABLE, StatusCode.DEADLINE_EXCEEDED, ]
)

with DaprClient(retry_policy=retry) as d:
    ...

或者对于 actors:

factory = ActorProxyFactory(retry_policy=RetryPolicy(max_attempts=3))
proxy = ActorProxy.create('DemoActor', ActorId('1'), DemoActorInterface, factory)

超时可以通过环境变量 DAPR_API_TIMEOUT_SECONDS 为所有调用设置。默认值为 60 秒。

注意:你可以单独控制服务调用的超时,方法是将 timeout 参数传递给 invoke_method 方法。

错误处理

最初,Dapr 中的错误遵循标准 gRPC 错误模型。然而,为了提供更详细和更有信息量的错误消息,在 1.13 版本中引入了一个增强的错误模型,该模型与 gRPC 更丰富的错误模型保持一致。作为响应,Python SDK 实现了 DaprGrpcError,这是一个自定义异常类,旨在改善开发人员体验。
值得注意的是,对于所有 gRPC 状态异常,使用 DaprGrpcError 的转换正在进行中。截至目前,并非 SDK 中的每个 API 调用都已更新以利用此自定义异常。我们正在积极进行此增强,并欢迎社区的贡献。

使用 Dapr python-SDK 时处理 DaprGrpcError 异常的示例:

try:
    d.save_state(store_name=storeName, key=key, value=value)
except DaprGrpcError as err:
    print(f'Status code: {err.code()}')
    print(f"Message: {err.message()}")
    print(f"Error code: {err.error_code()}")
    print(f"Error info(reason): {err.error_info.reason}")
    print(f"Resource info (resource type): {err.resource_info.resource_type}")
    print(f"Resource info (resource name): {err.resource_info.resource_name}")
    print(f"Bad request (field): {err.bad_request.field_violations[0].field}")
    print(f"Bad request (description): {err.bad_request.field_violations[0].description}")

构建块

Python SDK 允许你与所有 Dapr 构建块 进行交互。

调用服务

Dapr Python SDK 提供了一个简单的 API,可以通过 HTTP 或 gRPC(已弃用)调用服务。可以通过设置 DAPR_API_METHOD_INVOCATION_PROTOCOL 环境变量来选择协议,未设置时默认为 HTTP。Dapr 中的 GRPC 服务调用已弃用,建议使用 GRPC 代理作为替代。

from dapr.clients import DaprClient

with DaprClient() as d:
    # 调用方法(gRPC 或 HTTP GET)    
    resp = d.invoke_method('service-to-invoke', 'method-to-invoke', data='{"message":"Hello World"}')

    # 对于其他 HTTP 动词,必须指定动词
    # 调用 'POST' 方法(仅 HTTP)    
    resp = d.invoke_method('service-to-invoke', 'method-to-invoke', data='{"id":"100", "FirstName":"Value", "LastName":"Value"}', http_verb='post')

HTTP api 调用的基本端点在 DAPR_HTTP_ENDPOINT 环境变量中指定。 如果未设置此变量,端点值将从 DAPR_RUNTIME_HOSTDAPR_HTTP_PORT 变量派生,其默认值分别为 127.0.0.13500

gRPC 调用的基本端点是用于客户端初始化的端点(上文解释)。

保存和获取应用程序状态

from dapr.clients import DaprClient

with DaprClient() as d:
    # 保存状态
    d.save_state(store_name="statestore", key="key1", value="value1")

    # 获取状态
    data = d.get_state(store_name="statestore", key="key1").data

    # 删除状态
    d.delete_state(store_name="statestore", key="key1")

查询应用程序状态(Alpha)

    from dapr import DaprClient

    query = '''
    {
        "filter": {
            "EQ": { "state": "CA" }
        },
        "sort": [
            {
                "key": "person.id",
                "order": "DESC"
            }
        ]
    }
    '''

    with DaprClient() as d:
        resp = d.query_state(
            store_name='state_store',
            query=query,
            states_metadata={"metakey": "metavalue"},  # 可选
        )

发布和订阅

发布消息

from dapr.clients import DaprClient

with DaprClient() as d:
    resp = d.publish_event(pubsub_name='pubsub', topic_name='TOPIC_A', data='{"message":"Hello World"}')

发送带有 json 有效负载的 CloudEvents 消息:

from dapr.clients import DaprClient
import json

with DaprClient() as d:
    cloud_event = {
        'specversion': '1.0',
        'type': 'com.example.event',
        'source': 'my-service',
        'id': 'myid',
        'data': {'id': 1, 'message': 'hello world'},
        'datacontenttype': 'application/json',
    }

    # 将数据内容类型设置为 'application/cloudevents+json'
    resp = d.publish_event(
        pubsub_name='pubsub',
        topic_name='TOPIC_CE',
        data=json.dumps(cloud_event),
        data_content_type='application/cloudevents+json',
    )

发布带有纯文本有效负载的 CloudEvents 消息:

from dapr.clients import DaprClient
import json

with DaprClient() as d:
    cloud_event = {
        'specversion': '1.0',
        'type': 'com.example.event',
        'source': 'my-service',
        'id': "myid",
        'data': 'hello world',
        'datacontenttype': 'text/plain',
    }

    # 将数据内容类型设置为 'application/cloudevents+json'
    resp = d.publish_event(
        pubsub_name='pubsub',
        topic_name='TOPIC_CE',
        data=json.dumps(cloud_event),
        data_content_type='application/cloudevents+json',
    )

订阅消息

from cloudevents.sdk.event import v1
from dapr.ext.grpc import App
import json

app = App()

# 主题的默认订阅
@app.subscribe(pubsub_name='pubsub', topic='TOPIC_A')
def mytopic(event: v1.Event) -> None:
    data = json.loads(event.Data())
    print(f'Received: id={data["id"]}, message="{data ["message"]}"' 
          ' content_type="{event.content_type}"',flush=True)

# 使用发布/订阅路由的特定处理程序
@app.subscribe(pubsub_name='pubsub', topic='TOPIC_A',
               rule=Rule("event.type == \"important\"", 1))
def mytopic_important(event: v1.Event) -> None:
    data = json.loads(event.Data())
    print(f'Received: id={data["id"]}, message="{data ["message"]}"' 
          ' content_type="{event.content_type}"',flush=True)

流式消息订阅

你可以使用 subscribesubscribe_handler 方法创建对 PubSub 主题的流式订阅。

subscribe 方法返回一个可迭代的 Subscription 对象,允许你使用 for 循环(例如 for message in subscription)或调用 next_message 方法从流中拉取消息。这将在等待消息时阻塞主线程。 完成后,你应该调用 close 方法来终止订阅并停止接收消息。

subscribe_with_handler 方法接受一个回调函数,该函数对从流中接收到的每条消息执行。 它在单独的线程中运行,因此不会阻塞主线程。回调应返回一个 TopicEventResponse(例如 TopicEventResponse('success')),指示消息是成功处理、应该重试还是应该丢弃。该方法将根据返回的状态自动管理消息确认。对 subscribe_with_handler 方法的调用返回一个关闭函数,完成后应调用该函数以终止订阅。

以下是使用 subscribe 方法的示例:

import time

from dapr.clients import DaprClient
from dapr.clients.grpc.subscription import StreamInactiveError, StreamCancelledError

counter = 0


def process_message(message):
    global counter
    counter += 1
    # 在此处处理消息
    print(f'Processing message: {message.data()} from {message.topic()}...')
    return 'success'


def main():
    with DaprClient() as client:
        global counter

        subscription = client.subscribe(
            pubsub_name='pubsub', topic='TOPIC_A', dead_letter_topic='TOPIC_A_DEAD'
        )

        try:
            for message in subscription:
                if message is None:
                    print('No message received. The stream might have been cancelled.')
                    continue

                try:
                    response_status = process_message(message)

                    if response_status == 'success':
                        subscription.respond_success(message)
                    elif response_status == 'retry':
                        subscription.respond_retry(message)
                    elif response_status == 'drop':
                        subscription.respond_drop(message)

                    if counter >= 5:
                        break
                except StreamInactiveError:
                    print('Stream is inactive. Retrying...')
                    time.sleep(1)
                    continue
                except StreamCancelledError:
                    print('Stream was cancelled')
                    break
                except Exception as e:
                    print(f'Error occurred during message processing: {e}')

        finally:
            print('Closing subscription...')
            subscription.close()


if __name__ == '__main__':
    main()

以下是使用 subscribe_with_handler 方法的示例:

import time

from dapr.clients import DaprClient
from dapr.clients.grpc._response import TopicEventResponse

counter = 0


def process_message(message):
    # 在此处处理消息
    global counter
    counter += 1
    print(f'Processing message: {message.data()} from {message.topic()}...')
    return TopicEventResponse('success')


def main():
    with (DaprClient() as client):
        # 这将启动一个新线程来监听消息
        # 并在 `process_message` 函数中处理它们
        close_fn = client.subscribe_with_handler(
            pubsub_name='pubsub', topic='TOPIC_A', handler_fn=process_message,
            dead_letter_topic='TOPIC_A_DEAD'
        )

        while counter < 5:
            time.sleep(1)

        print("Closing subscription...")
        close_fn()


if __name__ == '__main__':
    main()

对话(Alpha)

从 1.15 版本开始,Dapr 为开发者提供了通过对话 API与大型语言模型(LLM)安全可靠地交互的能力。

from dapr.clients import DaprClient
from dapr.clients.grpc.conversation import ConversationInput

with DaprClient() as d:
    inputs = [
        ConversationInput(content="What's Dapr?", role='user', scrub_pii=True),
        ConversationInput(content='Give a brief overview.', role='user', scrub_pii=True),
    ]

    metadata = {
        'model': 'foo',
        'key': 'authKey',
        'cacheTTL': '10m',
    }

    response = d.converse_alpha1(
        name='echo', inputs=inputs, temperature=0.7, context_id='chat-123', metadata=metadata
    )

    for output in response.outputs:
        print(f'Result: {output.result}')

与输出绑定交互

from dapr.clients import DaprClient

with DaprClient() as d:
    resp = d.invoke_binding(binding_name='kafkaBinding', operation='create', data='{"message":"Hello World"}')

检索密钥

from dapr.clients import DaprClient

with DaprClient() as d:
    resp = d.get_secret(store_name='localsecretstore', key='secretKey')

配置

获取配置

from dapr.clients import DaprClient

with DaprClient() as d:
    # 获取配置
    configuration = d.get_configuration(store_name='configurationstore', keys=['orderId'], config_metadata={})

订阅配置

import asyncio
from time import sleep
from dapr.clients import DaprClient

async def executeConfiguration():
    with DaprClient() as d:
        storeName = 'configurationstore'

        key = 'orderId'

        # 等待边车在 20 秒内启动。
        d.wait(20)

        # 按键订阅配置。
        configuration = await d.subscribe_configuration(store_name=storeName, keys=[key], config_metadata={})
        while True:
            if configuration != None:
                items = configuration.get_items()
                for key, item in items:
                    print(f"Subscribe key={key} value={item.value} version={item.version}", flush=True)
            else:
                print("Nothing yet")
        sleep(5)

asyncio.run(executeConfiguration())

分布式锁

from dapr.clients import DaprClient

def main():
    # 锁参数
    store_name = 'lockstore'  # 在 components/lockstore.yaml 中定义
    resource_id = 'example-lock-resource'
    client_id = 'example-client-id'
    expiry_in_seconds = 60

    with DaprClient() as dapr:
        print('Will try to acquire a lock from lock store named [%s]' % store_name)
        print('The lock is for a resource named [%s]' % resource_id)
        print('The client identifier is [%s]' % client_id)
        print('The lock will will expire in %s seconds.' % expiry_in_seconds)

        with dapr.try_lock(store_name, resource_id, client_id, expiry_in_seconds) as lock_result:
            assert lock_result.success, 'Failed to acquire the lock. Aborting.'
            print('Lock acquired successfully!!!')

        # 此时锁已释放 - 通过 `with` 子句的魔法 ;)
        unlock_result = dapr.unlock(store_name, resource_id, client_id)
        print('We already released the lock so unlocking will not work.')
        print('We tried to unlock it anyway and got back [%s]' % unlock_result.status)

加密

from dapr.clients import DaprClient

message = 'The secret is "passw0rd"'

def main():
    with DaprClient() as d:
        resp = d.encrypt(
            data=message.encode(),
            options=EncryptOptions(
                component_name='crypto-localstorage',
                key_name='rsa-private-key.pem',
                key_wrap_algorithm='RSA',
            ),
        )
        encrypt_bytes = resp.read()

        resp = d.decrypt(
            data=encrypt_bytes,
            options=DecryptOptions(
                component_name='crypto-localstorage',
                key_name='rsa-private-key.pem',
            ),
        )
        decrypt_bytes = resp.read()

        print(decrypt_bytes.decode())  # The secret is "passw0rd"

相关链接

Python SDK 示例

2 - Dapr Actor Python SDK 入门

如何快速上手使用 Dapr Python SDK

Dapr actor 包使您能够从 Python 应用程序与 Dapr virtual actors 交互。

前置条件

Actor 接口

接口定义了 actor 实现与调用 actor 的客户端之间共享的 actor 契约。由于客户端可能依赖该接口,因此通常将其定义在与 actor 实现分离的程序集中。

from dapr.actor import ActorInterface, actormethod

class DemoActorInterface(ActorInterface):
    @actormethod(name="GetMyData")
    async def get_my_data(self) -> object:
        ...

Actor 服务

actor 服务托管 virtual actor。它由一个派生自基类型 Actor 并实现 actor 接口中定义的接口的类来实现。

可以使用以下 Dapr actor 扩展之一创建 actor:

Actor 客户端

actor 客户端包含 actor 客户端的实现,该实现调用 actor 接口中定义的 actor 方法。

import asyncio

from dapr.actor import ActorProxy, ActorId
from demo_actor_interface import DemoActorInterface

async def main():
    # 创建代理客户端
    proxy = ActorProxy.create('DemoActor', ActorId('1'), DemoActorInterface)

    # 在客户端上调用方法
    resp = await proxy.GetMyData()

示例

请访问此页面查看可运行的 actor 示例。

Mock Actor 测试

Dapr Python SDK 提供了创建 mock actor 的功能,用于对 actor 方法进行单元测试,并查看它们如何与 actor 状态交互。

示例用法

from dapr.actor.runtime.mock_actor import create_mock_actor

class MyActor(Actor, MyActorInterface):
    async def save_state(self, data) -> None:
        await self._state_manager.set_state('mystate', data)
        await self._state_manager.save_state()

mock_actor = create_mock_actor(MyActor, "id")

await mock_actor.save_state(5)
assert mockactor._state_manager._mock_state['mystate'] == 5 #True

Mock actor 通过将 actor 类和 actor ID(字符串)传递给 create_mock_actor 函数来创建。此函数返回一个 actor 实例,其中许多内部方法已被覆盖。Mock actor 不与 Dapr 交互来执行保存状态或管理定时器等任务,而是使用内存状态来模拟这些行为。

可以通过以下变量访问此状态:

重要提示:由于下文详细讨论的类型提示问题,这些变量对类型提示器/linter 等工具不可见,它们会认为这些是无效变量。您需要使用 #type: ignore 来满足任何此类系统的要求。

  • _state_manager._mock_state()
    一个 [str, object] 字典,其中存储了所有 actor 状态。通过 _state_manager.save_state(key, value) 或任何其他 statemanager 方法保存的任何变量都作为该键值对存储在字典中。通过 try_get_state 或任何其他 statemanager 方法加载的任何值都从此字典中获取。

  • _state_manager._mock_timers()
    一个 [str, ActorTimerData] 字典,其中保存活动的 actor 定时器。任何会添加或移除定时器的 actor 方法都会从此字典中添加或弹出相应的 ActorTimerData 对象。

  • _state_manager._mock_reminders()
    一个 [str, ActorReminderData] 字典,其中保存活动的 actor 提醒。任何会添加或移除定时器的 actor 方法都会从此字典中添加或弹出相应的 ActorReminderData 对象。

注意:定时器和提醒永远不会实际触发。这些字典的存在仅为了测试应该添加或移除定时器/提醒的方法。如果您需要测试它们应该激活的回调,您应该使用适当的值直接调用它们:

result = await mock_actor.recieve_reminder(name, state, due_time, period, _ttl)
# 直接测试结果,或通过查询 `_state_manager._mock_state` 测试副作用(如更改状态)

用法和限制

为了允许更细粒度的控制,_on_activate 方法不会像 Dapr 初始化新的 Actor 实例时那样自动调用。您应该在测试中根据需要手动调用它。

Mock actor 系统当前的一个限制是它不调用 _on_pre_actor_method_on_post_actor_method 方法。您始终可以手动调用这些方法作为测试的一部分。

__init__register_timerunregister_timerregister_reminderunregister_reminder 方法都会被 MockActor 类覆盖,该类通过 create_mock_actor 作为 mixin 应用。如果您的 actor 本身覆盖了这些方法,这些修改本身将被覆盖,actor 可能不会按您的预期运行。

注意:__init__ 是一个特殊情况,您应该将其定义为

    def __init__(self, ctx, actor_id):
        super().__init__(ctx, actor_id)

Mock actor 可以正常工作,但如果您在 __init__ 中添加了任何额外逻辑,它将被覆盖。值得注意的是,在初始化时应用逻辑的正确方法是通过 _on_activate(也可以与 mock actor 安全使用),而不是 __init__

如果您有一个确实覆盖了默认 Dapr actor 方法的 actor,您可以创建 MockActor 类(来自 MockActor.py)的自定义子类,该类实现您拥有的任何自定义逻辑,同时与 _mock_state_mock_timers_mock_reminders 正常交互,然后通过您自己定义的 create_mock_actor 函数将该自定义类作为 mixin 应用。

Actor _runtime_ctx 变量设置为 None。所有正常的 actor 方法都已覆盖,不会调用它,但如果您的代码本身直接与 _runtime_ctx 交互,测试可能会失败。

Actor _state_manager 被 MockStateManager 实例覆盖。它具有与基本 ActorStateManager 相同的所有方法和功能,除了使用各种 _mock 变量而不是 _runtime_ctx 来存储数据。如果您的代码实现了自己的自定义状态管理器,它将被覆盖,测试可能会失败。

类型提示

由于 Python 缺乏用于类型提示类型交集的统一方法(参见:python/typing #213),类型提示不幸地不适用于 Mock Actor。返回类型被类型提示为"Actor 子类 T 的实例",而实际上应该被类型提示为"MockActor 子类 T 的实例"或"类型交集 [Actor 子类 T, MockActor] 的实例"(值得注意的是,MockActor 本身是 Actor 的子类)。

这意味着,例如,如果您在代码编辑器中将鼠标悬停在 mockactor._state_manager 上,它将显示为 ActorStateManager 的实例(而不是 MockStateManager),各种 IDE 辅助功能(如 VSCode 的 Go to Definition,它会将您带到 ActorStateManager 的定义而不是 MockStateManager)将无法正常工作。

目前,这个问题无法修复,因此仅仅需要注意它,因为它可能会引起混淆。如果将来能够准确地为这种情况进行类型提示,欢迎随时打开关于实现它的问题。

3 - Dapr Python SDK 扩展

Python SDK,用于开发 Dapr 应用程序

3.1 - Dapr Python gRPC 服务扩展入门

如何快速上手 Dapr Python gRPC 扩展

Dapr Python SDK 提供了一个内置的 gRPC 服务器扩展 dapr.ext.grpc,用于创建 Dapr 服务。

安装

您可以通过以下命令下载并安装 Dapr gRPC 服务器扩展:

pip install dapr-ext-grpc
pip3 install dapr-ext-grpc-dev

示例

App 对象可用于创建服务器。

监听服务调用请求

InvokeMethodReqestInvokeMethodResponse 对象可用于处理传入请求。

一个简单的监听并响应请求的服务如下所示:

from dapr.ext.grpc import App, InvokeMethodRequest, InvokeMethodResponse

app = App()

@app.method(name='my-method')
def mymethod(request: InvokeMethodRequest) -> InvokeMethodResponse:
    print(request.metadata, flush=True)
    print(request.text(), flush=True)

    return InvokeMethodResponse(b'INVOKE_RECEIVED', "text/plain; charset=UTF-8")

app.run(50051)

完整示例可在此处找到。

订阅主题

在订阅主题时,您可以指示 Dapr 已接受传递的事件,还是应该丢弃该事件或稍后重试。

from typing import Optional
from cloudevents.sdk.event import v1
from dapr.ext.grpc import App
from dapr.clients.grpc._response import TopicEventResponse

app = App()

# 主题的默认订阅
@app.subscribe(pubsub_name='pubsub', topic='TOPIC_A')
def mytopic(event: v1.Event) -> Optional[TopicEventResponse]:
    print(event.Data(),flush=True)
    # 返回 None(或不显式返回)等效于
    # 返回 TopicEventResponse("success")。
    # 您也可以返回 TopicEventResponse("retry") 让 dapr 记录
    # 该消息并稍后重试传递,或返回 TopicEventResponse("drop")
    # 让其丢弃该消息
    return TopicEventResponse("success")

# 使用发布订阅路由的特定处理程序
@app.subscribe(pubsub_name='pubsub', topic='TOPIC_A',
               rule=Rule("event.type == \"important\"", 1))
def mytopic_important(event: v1.Event) -> None:
    print(event.Data(),flush=True)

# 禁用主题验证的处理程序
@app.subscribe(pubsub_name='pubsub-mqtt', topic='topic/#', disable_topic_validation=True,)
def mytopic_wildcard(event: v1.Event) -> None:
    print(event.Data(),flush=True)

app.run(50051)

完整示例可在此处找到。

设置输入绑定触发器

from dapr.ext.grpc import App, BindingRequest

app = App()

@app.binding('kafkaBinding')
def binding(request: BindingRequest):
    print(request.text(), flush=True)

app.run(50051)

完整示例可在此处找到。

相关链接

3.2 - Dapr Python SDK 与 FastAPI 集成

如何使用 FastAPI 扩展创建 Dapr Python virtual actors 和发布订阅

Dapr Python SDK 通过 dapr-ext-fastapi 扩展提供与 FastAPI 的集成。

安装

您可以通过以下命令下载并安装 Dapr FastAPI 扩展:

pip install dapr-ext-fastapi
pip install dapr-ext-fastapi-dev

示例

订阅不同类型的事件

import uvicorn
from fastapi import Body, FastAPI
from dapr.ext.fastapi import DaprApp
from pydantic import BaseModel

class RawEventModel(BaseModel):
    body: str

class User(BaseModel):
    id: int
    name: str

class CloudEventModel(BaseModel):
    data: User
    datacontenttype: str
    id: str
    pubsubname: str
    source: str
    specversion: str
    topic: str
    traceid: str
    traceparent: str
    tracestate: str
    type: str    
    
    
app = FastAPI()
dapr_app = DaprApp(app)

# 允许处理任何结构的事件(最简单,但最不健壮)
# dapr publish --publish-app-id sample --topic any_topic --pubsub pubsub --data '{"id":"7", "desc": "good", "size":"small"}'
@dapr_app.subscribe(pubsub='pubsub', topic='any_topic')
def any_event_handler(event_data = Body()):
    print(event_data)    

# 为了健壮性,根据发布者是否使用 CloudEvents 选择以下方式之一

# 处理使用 CloudEvents 发送的事件
# dapr publish --publish-app-id sample --topic cloud_topic --pubsub pubsub --data '{"id":"7", "name":"Bob Jones"}'
@dapr_app.subscribe(pubsub='pubsub', topic='cloud_topic')
def cloud_event_handler(event_data: CloudEventModel):
    print(event_data)   

# 处理不使用 CloudEvents 发送的原始事件
# curl -X "POST" http://localhost:3500/v1.0/publish/pubsub/raw_topic?metadata.rawPayload=true -H "Content-Type: application/json" -d '{"body": "345"}'
@dapr_app.subscribe(pubsub='pubsub', topic='raw_topic')
def raw_event_handler(event_data: RawEventModel):
    print(event_data)    

 

if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=30212)

创建 actor

from fastapi import FastAPI
from dapr.ext.fastapi import DaprActor
from demo_actor import DemoActor

app = FastAPI(title=f'{DemoActor.__name__}Service')

# 添加 Dapr Actor 扩展
actor = DaprActor(app)

@app.on_event("startup")
async def startup_event():
    # 注册 DemoActor
    await actor.register_actor(DemoActor)

@app.get("/GetMyData")
def get_my_data():
    return "{'message': 'myData'}"

3.3 - Dapr Python SDK 与 Flask 集成

如何使用 Flask 扩展创建 Dapr Python virtual actors

Dapr Python SDK 通过 flask-dapr 扩展提供与 Flask 的集成。

安装

你可以使用以下命令下载并安装 Dapr Flask 扩展:

pip install flask-dapr
pip install flask-dapr-dev

示例

from flask import Flask
from flask_dapr.actor import DaprActor

from dapr.conf import settings
from demo_actor import DemoActor

app = Flask(f'{DemoActor.__name__}Service')

# 启用 DaprActor Flask 扩展
actor = DaprActor(app)

# 注册 DemoActor
actor.register_actor(DemoActor)

# 设置方法路由
@app.route('/GetMyData', methods=['GET'])
def get_my_data():
    return {'message': 'myData'}, 200

# 运行应用
if __name__ == '__main__':
    app.run(port=settings.HTTP_APP_PORT)

3.4 - Dapr Python SDK 与 Dapr 工作流扩展集成

如何快速上手使用 Dapr 工作流扩展

Dapr Python SDK 提供了一个内置的 Dapr 工作流扩展 dapr.ext.workflow,用于创建 Dapr 服务。

安装

您可以通过以下命令下载并安装 Dapr 工作流扩展:

pip install dapr-ext-workflow
pip install dapr-ext-workflow-dev

示例

from time import sleep

import dapr.ext.workflow as wf


wfr = wf.WorkflowRuntime()


@wfr.workflow(name='random_workflow')
def task_chain_workflow(ctx: wf.DaprWorkflowContext, wf_input: int):
    try:
        result1 = yield ctx.call_activity(step1, input=wf_input)
        result2 = yield ctx.call_activity(step2, input=result1)
    except Exception as e:
        yield ctx.call_activity(error_handler, input=str(e))
        raise
    return [result1, result2]


@wfr.activity(name='step1')
def step1(ctx, activity_input):
    print(f'Step 1: Received input: {activity_input}.')
    # Do some work
    return activity_input + 1


@wfr.activity
def step2(ctx, activity_input):
    print(f'Step 2: Received input: {activity_input}.')
    # Do some work
    return activity_input * 2

@wfr.activity
def error_handler(ctx, error):
    print(f'Executing error handler: {error}.')
    # Do some compensating work


if __name__ == '__main__':
    wfr.start()
    sleep(10)  # wait for workflow runtime to start

    wf_client = wf.DaprWorkflowClient()
    instance_id = wf_client.schedule_new_workflow(workflow=task_chain_workflow, input=42)
    print(f'Workflow started. Instance ID: {instance_id}')
    state = wf_client.wait_for_workflow_completion(instance_id)
    print(f'Workflow completed! Status: {state.runtime_status}')

    wfr.shutdown()

后续步骤

Dapr 工作流 Python SDK 入门

3.4.1 - Dapr Workflow Python SDK 入门

如何使用 Dapr Python SDK 快速上手工作流

让我们创建一个 Dapr 工作流并通过控制台调用它。借助提供的工作流示例,你将:

  • 运行一个 Python 控制台应用程序,该程序演示包含活动、子工作流和外部事件的工作流编排
  • 了解如何处理重试、超时以及工作流状态管理
  • 使用 Python 工作流 SDK 来启动、暂停、恢复和清理工作流实例

本示例使用自托管模式下通过 dapr init 初始化的默认配置。

在 Python 示例项目中,simple.py 文件包含应用程序的设置,包括:

  • 工作流定义
  • 工作流活动定义
  • 工作流和工作流活动的注册

前置条件

设置环境

首先克隆 [Python SDK 仓库]。

git clone https://github.com/dapr/python-sdk.git

从 Python SDK 根目录导航到 Dapr Workflow 示例。

cd examples/workflow

运行以下命令,安装使用 Dapr Python SDK 运行此工作流示例所需的所有依赖。

pip3 install -r workflow/requirements.txt

在本地运行应用程序

要运行 Dapr 应用程序,你需要启动 Python 程序和一个 Dapr 边车。在终端中运行:

dapr run --app-id wf-simple-example --dapr-grpc-port 50001 --resources-path components -- python3 simple.py

注意: 由于 Windows 上未定义 Python3.exe,你可能需要使用 python simple.py 而不是 python3 simple.py

预期输出

- "== APP == Hi Counter!"
- "== APP == New counter value is: 1!"
- "== APP == New counter value is: 11!"
- "== APP == Retry count value is: 0!"
- "== APP == Retry count value is: 1! This print statement verifies retry"
- "== APP == Appending 1 to child_orchestrator_string!"
- "== APP == Appending a to child_orchestrator_string!"
- "== APP == Appending a to child_orchestrator_string!"
- "== APP == Appending 2 to child_orchestrator_string!"
- "== APP == Appending b to child_orchestrator_string!"
- "== APP == Appending b to child_orchestrator_string!"
- "== APP == Appending 3 to child_orchestrator_string!"
- "== APP == Appending c to child_orchestrator_string!"
- "== APP == Appending c to child_orchestrator_string!"
- "== APP == Get response from hello_world_wf after pause call: Suspended"
- "== APP == Get response from hello_world_wf after resume call: Running"
- "== APP == New counter value is: 111!"
- "== APP == New counter value is: 1111!"
- "== APP == Workflow completed! Result: "Completed"

发生了什么?

当你运行应用程序时,会演示几个关键的工作流功能:

  1. 工作流和活动注册:应用程序使用 Python 装饰器自动向运行时注册工作流和活动。这种基于装饰器的方法提供了一种简洁、声明式的方式来定义你的工作流组件:

    @wfr.workflow(name='hello_world_wf')
    def hello_world_wf(ctx: DaprWorkflowContext, wf_input):
        # Workflow definition...
    
    @wfr.activity(name='hello_act')
    def hello_act(ctx: WorkflowActivityContext, wf_input):
        # Activity definition...
    
  2. 运行时设置:应用程序初始化工作流运行时和客户端:

    wfr = WorkflowRuntime()
    wfr.start()
    wf_client = DaprWorkflowClient()
    
  3. 活动执行:工作流执行一系列活动来递增计数器:

    @wfr.workflow(name='hello_world_wf')
    def hello_world_wf(ctx: DaprWorkflowContext, wf_input):
        yield ctx.call_activity(hello_act, input=1)
        yield ctx.call_activity(hello_act, input=10)
    
  4. 重试逻辑:工作流演示了使用重试策略进行错误处理:

    retry_policy = RetryPolicy(
        first_retry_interval=timedelta(seconds=1),
        max_number_of_attempts=3,
        backoff_coefficient=2,
        max_retry_interval=timedelta(seconds=10),
        retry_timeout=timedelta(seconds=100),
    )
    yield ctx.call_activity(hello_retryable_act, retry_policy=retry_policy)
    
  5. 子工作流:子工作流使用自己的重试策略执行:

    yield ctx.call_child_workflow(child_retryable_wf, retry_policy=retry_policy)
    
  6. 外部事件处理:工作流等待一个带有超时的外部事件:

    event = ctx.wait_for_external_event(event_name)
    timeout = ctx.create_timer(timedelta(seconds=30))
    winner = yield when_any([event, timeout])
    
  7. 工作流生命周期管理:示例演示如何暂停和恢复工作流:

    wf_client.pause_workflow(instance_id=instance_id)
    metadata = wf_client.get_workflow_state(instance_id=instance_id)
    # ... check status ...
    wf_client.resume_workflow(instance_id=instance_id)
    
  8. 事件触发:恢复后,工作流触发一个事件:

    wf_client.raise_workflow_event(
        instance_id=instance_id,
        event_name=event_name,
        data=event_data
    )
    
  9. 完成与清理:最后,工作流等待完成并进行清理:

    state = wf_client.wait_for_workflow_completion(
        instance_id,
        timeout_in_seconds=30
    )
    wf_client.purge_workflow(instance_id=instance_id)
    

后续步骤

4 -

title: “Conversation API(Python)- 推荐用法” linkTitle: “对话” weight: 11000 type: docs description: 推荐在 Python 中配合或不配合工具使用 Dapr Conversation API 的模式,包括多轮流程与安全指引。

Dapr Conversation API 当前仍处于 alpha 阶段。本文给出在 Python SDK 中高效使用它的推荐最小模式:

  • 纯请求(无工具)
  • 携带工具的请求(将函数用作工具)
  • 带工具执行的多轮流程
  • 异步变体
  • 执行工具调用时的重要安全注意事项

前提条件

如需完整的端到端流程与提供程序配置,参见:

基础对话(无工具)

from dapr.clients import DaprClient
from dapr.clients.grpc import conversation

# 构建单轮 Alpha2 输入
user_msg = conversation.create_user_message("什么是 Dapr?")
alpha2_input = conversation.ConversationInputAlpha2(messages=[user_msg])

with DaprClient() as client:
    resp = client.converse_alpha2(
        name="echo",  # 替换为你的 LLM 组件名称
        inputs=[alpha2_input],
        temperature=1,
    )

    for msg in resp.to_assistant_messages():
        if msg.of_assistant.content:
            print(msg.of_assistant.content[0].text)

要点:

  • 使用 conversation.create_user_message 构造消息。
  • 将其包进 ConversationInputAlpha2(messages=[...]),再传给 converse_alpha2
  • 使用 response.to_assistant_messages() 遍历 assistant 输出。

工具:基于装饰器(推荐)

基于装饰器的工具方式更整洁、也更顺手。定义函数时应写清晰的类型提示与详细 docstring;这对 LLM 理解何时以及如何调用该工具很重要; 然后使用 @conversation.tool 进行装饰。注册后的工具可以传给 LLM,并通过工具调用来执行。

from dapr.clients import DaprClient
from dapr.clients.grpc import conversation

@conversation.tool
def get_weather(location: str, unit: str = 'fahrenheit') -> str:
    """获取指定地点的当前天气。"""
    # 请替换为真实实现
    return f"Weather in {location} (unit={unit})"

user_msg = conversation.create_user_message("巴黎的天气怎么样?")
alpha2_input = conversation.ConversationInputAlpha2(messages=[user_msg])

with DaprClient() as client:
    response = client.converse_alpha2(
        name="openai",  # 你的 LLM 组件
        inputs=[alpha2_input],
        tools=conversation.get_registered_tools(),  # 由 @conversation.tool 注册的工具
        tool_choice='auto',
        temperature=1,
    )

    # 检查 assistant 消息,包括其中的工具调用
    for msg in response.to_assistant_messages():
        if msg.of_assistant.tool_calls:
            for tc in msg.of_assistant.tool_calls:
                print(f"Tool call: {tc.function.name} args={tc.function.arguments}")
        elif msg.of_assistant.content:
            print(msg.of_assistant.content[0].text)

说明:

  • 使用 conversation.get_registered_tools() 收集所有通过 @conversation.tool 装饰的函数。
  • 绑定器会依据函数签名校验并强制转换参数。类型注解要准确。

使用工具的最小多轮流程

这是处理带工具对话时最常用的循环:

from dapr.clients import DaprClient
from dapr.clients.grpc import conversation

@conversation.tool
def get_weather(location: str, unit: str = 'fahrenheit') -> str:
    return f"Weather in {location} (unit={unit})"

history: list[conversation.ConversationMessage] = [
    conversation.create_user_message("旧金山的天气怎么样?")]

with DaprClient() as client:
    # 第 1 轮
    resp1 = client.converse_alpha2(
        name="openai",
        inputs=[conversation.ConversationInputAlpha2(messages=history)],
        tools=conversation.get_registered_tools(),
        tool_choice='auto',
        temperature=1,
    )

    # 追加 assistant 消息;执行工具调用;再追加工具结果
    for msg in resp1.to_assistant_messages():
        history.append(msg)
        for tc in msg.of_assistant.tool_calls:
            # 重要:生产环境中必须校验输入并施加防护
            tool_output = conversation.execute_registered_tool(
                tc.function.name, tc.function.arguments
            )
            history.append(
                conversation.create_tool_message(
                    tool_id=tc.id, name=tc.function.name, content=str(tool_output)
                )
            )

    # 第 2 轮(LLM 会看到工具结果)
    history.append(conversation.create_user_message("我需要带伞吗?"))
    resp2 = client.converse_alpha2(
        name="openai",
        inputs=[conversation.ConversationInputAlpha2(messages=history)],
        tools=conversation.get_registered_tools(),
        temperature=1,
    )

    for msg in resp2.to_assistant_messages():
        history.append(msg)
        if not msg.of_assistant.tool_calls and msg.of_assistant.content:
            print(msg.of_assistant.content[0].text)

提示:

  • 始终把 assistant 消息追加到 history。
  • 执行每个工具调用时,都要先校验,再把工具输出作为 tool message 追加进去。
  • 下一轮请求会带上这些工具结果,LLM 才能基于它们继续推理。

将函数用作工具:其他方式

当装饰器方式不方便时,还有两种选择。

A)从带类型的函数自动生成 schema:

from enum import Enum
from dapr.clients.grpc import conversation

class Units(Enum):
    CELSIUS = 'celsius'
    FAHRENHEIT = 'fahrenheit'

def get_weather(location: str, unit: Units = Units.FAHRENHEIT) -> str:
    return f"Weather in {location}"

fn = conversation.ConversationToolsFunction.from_function(get_weather)
weather_tool = conversation.ConversationTools(function=fn)

B)手写 JSON Schema(兜底方式):

from dapr.clients.grpc import conversation

fn = conversation.ConversationToolsFunction(
    name='get_weather',
    description='Get current weather',
    parameters={
        'type': 'object',
        'properties': {
            'location': {'type': 'string'},
            'unit': {'type': 'string', 'enum': ['celsius', 'fahrenheit']},
        },
        'required': ['location'],
    },
)
weather_tool = conversation.ConversationTools(function=fn)

异步变体

按需使用异步客户端与异步工具执行辅助方法。

import asyncio
from dapr.aio.clients import DaprClient as AsyncDaprClient
from dapr.clients.grpc import conversation

@conversation.tool
def get_time() -> str:
    return '2025-01-01T12:00:00Z'

async def main():
    async with AsyncDaprClient() as client:
        msg = conversation.create_user_message('现在几点?')
        inp = conversation.ConversationInputAlpha2(messages=[msg])
        resp = await client.converse_alpha2(
            name='openai', inputs=[inp], tools=conversation.get_registered_tools()
        )
        for m in resp.to_assistant_messages():
            if m.of_assistant.content:
                print(m.of_assistant.content[0].text)

asyncio.run(main())

如果你需要异步执行工具(例如网络 I/O),请实现异步函数,并结合超时参数使用 conversation.execute_registered_tool_async

安全与验证(必读)

LLM 可能会建议调用工具。所有由模型提供的参数都必须视为不可信输入。

建议:

  • 仅将可信函数注册为工具。为清晰性与自动 schema 生成,优先使用 @conversation.tool 装饰器。
  • 使用精确的类型注解和 docstring。SDK 会把函数签名转换成 JSON schema,并在参数绑定时执行类型强制转换,同时拒绝意外或无效字段。
  • 对可能产生副作用的工具(文件系统、网络、子进程)添加防护。可考虑白名单、沙箱和配额限制。
  • 执行前校验参数。例如清理文件路径,或限制 URL / 域名范围。
  • 考虑超时和并发控制。对于异步工具,可向 execute_registered_tool_async(..., timeout=...) 传入超时。
  • 记录并监控工具使用。默认拒绝:若校验失败,就不要执行工具,并以安全方式告知用户。

另请参阅 dapr/clients/grpc/conversation.py 中的内联说明(如 tool()ConversationToolsexecute_registered_tool),了解参数绑定与错误处理细节。

关键辅助方法(速查)

本节汇总了示例中使用到的 dapr.clients.grpc.conversation 辅助工具。

  • create_user_message(text: str) -> ConversationMessage

    • 为 Alpha2 构建 user 角色消息。可用于 history 列表。
    • 示例:history.append(conversation.create_user_message("Hello"))
  • create_system_message(text: str) -> ConversationMessage

    • 构建 system 消息,用于约束 assistant 的行为。
    • 示例:history = [conversation.create_system_message("You are a concise assistant.")]
  • create_assistant_message(text: str) -> ConversationMessage

    • 适合在测试或受控流程中注入 assistant 文本。
  • create_tool_message(tool_id: str, name: str, content: Any) -> ConversationMessage

    • 将工具输出转换为下一轮可供 LLM 读取的 tool message。
    • content 可以是任意对象;SDK 会安全地将其转为字符串。
    • 示例:history.append(conversation.create_tool_message(tool_id=tc.id, name=tc.function.name, content=conversation.execute_registered_tool(tc.function.name, tc.function.arguments)))
  • get_registered_tools() -> list[ConversationTools]

    • 返回当前进程内注册表中的全部工具。
    • 包括以下方式创建的工具:
      • @conversation.tool 装饰器(默认自动注册),以及
      • ConversationToolsFunction.from_functionregister=True(默认)。
    • converse_alpha2(..., tools=...) 中传入该列表。
  • register_tool(name: str, t: ConversationTools) / unregister_tool(name: str)

    • 手动管理工具注册表(例如高级场景、测试、清理)。
    • 名称必须唯一;在长生命周期进程中,记得注销以避免冲突。
  • execute_registered_tool(name: str, params: Mapping|Sequence|str|None) -> Any

    • 按名称同步执行已注册工具。
    • params 可接受 kwargs(mapping)、args(sequence)、JSON 字符串或 None。若传入 JSON 字符串(LLM 常见返回形式),SDK 会自动解析。
    • 参数会依据函数签名或 schema 进行校验与强制转换;多余或无效字段会直接报错。
    • 安全性:params 必须视为不可信输入;对副作用操作要加防护。
  • execute_registered_tool_async(name: str, params: Mapping|Sequence|str|None, *, timeout: float|None=None) -> Any

    • 异步版本。支持超时;对 I/O 密集型工具尤其建议使用。
    • 适用于异步工具,或配合 aio 客户端使用。
  • ConversationToolsFunction.from_function(func: Callable, register: bool = True) -> ConversationToolsFunction

    • 从带类型的 Python 函数(注解 + 可选 docstring)推导 JSON schema,并可选择直接注册为工具。
    • 常见用法:spec = conversation.ConversationToolsFunction.from_function(my_func);随后可依赖自动注册,也可用 ConversationTools(function=spec) 包装后调用 register_tool(spec.name, tool),或直接把 [tool] 传给 tools=
  • ConversationResponseAlpha2.to_assistant_messages() -> list[ConversationMessage]

    • 便捷方法,用于把响应输出转换为 assistant 的 ConversationMessage 对象,便于直接追加到 history(包括存在的 tool_calls)。

提示:@conversation.tool 装饰器是创建工具最简单的方式。它会根据函数自动生成 schema,支持可选的命名空间或名称覆盖,并自动注册工具(若想延后注册,可设置 register=False)。