Skip to main content

Dify对接钉钉开发者后台添加的机器人

前面给你的版本偏“概念代码”,这里改成工程化 services 分层结构

目标:

  • 钉钉 开发者后台机器人

  • Stream 模式

  • Python dingtalk-stream

  • FastAPI 项目结构

  • Dify Chat API Streaming

  • Service 层隔离

  • 后续可扩展 Redis、知识库、权限

dingtalk-stream SDK 本身就是针对钉钉 Stream 模式的机器人收消息、事件回调设计的。当前 Python 包版本可直接通过 pip 安装。(PyPI)


一、项目结构

dify-dingtalk-bot/

├── app/
│
│   ├── main.py
│   │
│   ├── config/
│   │   └── settings.py
│   │
│   ├── services/
│   │   │
│   │   ├── dify_service.py
│   │   │
│   │   ├── dingtalk_service.py
│   │   │
│   │   └── chat_service.py
│   │
│   ├── handlers/
│   │   └── robot_handler.py
│   │
│   └── utils/
│       └── logger.py
│
├── requirements.txt
└── .env

二、依赖

requirements.txt

fastapi
uvicorn
python-dotenv
requests
dingtalk-stream

安装:

pip install -r requirements.txt

三、配置

.env

# 钉钉Stream
DING_CLIENT_ID=dingxxxx
DING_CLIENT_SECRET=xxxx


# Dify
DIFY_API_KEY=app-xxxx
DIFY_API_URL=https://api.dify.ai/v1/chat-messages


# 服务
APP_NAME=dify-dingtalk-bot

四、配置服务

app/config/settings.py

import os
from dotenv import load_dotenv


load_dotenv()


class Settings:


    DING_CLIENT_ID = os.getenv(
        "DING_CLIENT_ID"
    )


    DING_CLIENT_SECRET = os.getenv(
        "DING_CLIENT_SECRET"
    )


    DIFY_API_KEY=os.getenv(
        "DIFY_API_KEY"
    )


    DIFY_API_URL=os.getenv(
        "DIFY_API_URL"
    )


settings=Settings()

五、Dify Service

负责:

  • 调用 Dify

  • 消费 SSE

  • 返回完整答案

app/services/dify_service.py

import requests
import json

from app.config.settings import settings



class DifyService:


    def __init__(self):

        self.url = (
            settings.DIFY_API_URL
        )

        self.key = (
            settings.DIFY_API_KEY
        )



    def chat(
        self,
        question:str,
        user:str
    ):


        headers={

            "Authorization":
            f"Bearer {self.key}",

            "Content-Type":
            "application/json"
        }


        payload={

            "inputs":{},

            "query":
            question,


            "response_mode":
            "streaming",


            "user":
            user

        }



        response=requests.post(

            self.url,

            headers=headers,

            json=payload,

            stream=True,

            timeout=120

        )


        answer=""


        for line in response.iter_lines():


            if not line:
                continue



            line=line.decode(
                "utf-8"
            )



            if not line.startswith(
                "data:"
            ):
                continue



            data=json.loads(
                line[5:]
            )


            event=data.get(
                "event"
            )


            if event=="message":

                answer += data.get(
                    "answer",
                    ""
                )


            if event=="message_end":

                break



        return answer

六、聊天业务 Service

以后加:

  • 用户权限

  • Redis

  • 会话ID

  • 日志

都放这里。

app/services/chat_service.py

from app.services.dify_service import DifyService



class ChatService:


    def __init__(self):

        self.dify=DifyService()



    def ask(
        self,
        message,
        user
    ):


        return self.dify.chat(

            question=message,

            user=user

        )

七、钉钉发送 Service

app/services/dingtalk_service.py

import dingtalk_stream



class DingTalkService:



    def reply(
        self,
        handler,
        text,
        message
    ):


        handler.reply_text(

            text,

            message

        )

八、机器人 Handler

这是核心。

app/handlers/robot_handler.py

import logging

import dingtalk_stream

from dingtalk_stream import AckMessage


from app.services.chat_service import ChatService
from app.services.dingtalk_service import DingTalkService



class RobotHandler(
    dingtalk_stream.ChatbotHandler
):


    def __init__(self):

        super().__init__()

        self.chat_service=(
            ChatService()
        )


        self.ding_service=(
            DingTalkService()
        )



    async def process(
        self,
        callback:
        dingtalk_stream.CallbackMessage
    ):


        incoming = (
            dingtalk_stream
            .ChatbotMessage
            .from_dict(
                callback.data
            )
        )


        question=(

            incoming
            .text
            .content
            .strip()

        )


        user=(

            incoming
            .sender_staff_id

        )


        logging.info(
            f"user={user}, msg={question}"
        )



        try:


            answer=(

                self.chat_service.ask(

                    question,

                    user

                )

            )


            self.ding_service.reply(

                self,

                answer,

                incoming

            )


        except Exception as e:


            self.reply_text(

                "服务异常:"+str(e),

                incoming

            )



        return (
            AckMessage.STATUS_OK,
            "OK"
        )

九、启动 Stream Client

app/main.py

import logging

import dingtalk_stream


from app.config.settings import settings

from app.handlers.robot_handler import RobotHandler



def start():


    logging.basicConfig(

        level=logging.INFO

    )



    credential=(

        dingtalk_stream
        .Credential(

            settings.DING_CLIENT_ID,

            settings.DING_CLIENT_SECRET

        )

    )



    client=(

        dingtalk_stream
        .DingTalkStreamClient(

            credential

        )

    )



    client.register_callback_handler(

        dingtalk_stream
        .chatbot
        .ChatbotMessage
        .TOPIC,


        RobotHandler()

    )



    client.start_forever()



if __name__=="__main__":

    start()

十、启动

python -m app.main

正常日志:

INFO endpoint is wss://xxxx
INFO websocket connected

钉钉开发者后台机器人上线后:

群里:

@机器人 公司的年假政策是什么?

流程:

ChatbotMessage
       |
       ↓
RobotHandler
       |
       ↓
ChatService
       |
       ↓
DifyService
       |
       ↓
Dify SSE
       |
       ↓
reply_text()
       |
       ↓
钉钉群

十一、下一步建议改造(生产版)

你这个架构建议继续升级:

1. 增加 Redis 会话

现在:

用户
 |
Dify
 |
新会话

升级:

staff_id
   |
Redis
   |
conversation_id
   |
Dify上下文

2. 增加真正 Stream 卡片回复

当前:

Dify:
chunk1
chunk2
chunk3


钉钉:
一次回复

升级:

Dify chunk
    |
    |
钉钉AI卡片stream
    |
    |
逐字刷新

可以做到类似 ChatGPT 的效果。


3. 推荐生产目录

大型项目:

services

├── dingtalk
│    ├── stream_client.py
│    ├── message.py
│    └── card.py
│
├── dify
│    ├── client.py
│    ├── stream.py
│    └── workflow.py
│
├── memory
│    └── redis.py
│
└── security
     └── permission.py

这个结构可以直接扩展成企业级 钉钉 AI 助手 + Dify RAG 知识库平台