You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Eclipse Ditto WebSocket更新数字孪生后Kafka主题未接收事件问题

Eclipse Ditto WebSocket更新孪生后Kafka主题收不到数据排查

问题复现环境与现象

  • 需求目标:通过Python编写的WebSocket客户端发送指令更新数字孪生的features属性,同时从Apache Kafka主题读取更新后的属性值
  • 实际运行表现:
    • WebSocket发送消息后Ditto可正常完成数字孪生更新
    • 第二个直连Ditto的WebSocket客户端可以正常接收到更新后的features属性
    • 配置好的Kafka目标主题无法接收到对应更新数据
    • 特殊场景:如果通过Kafka源连接指定主题触发数字孪生更新,目标Kafka主题可以正常收到更新后的属性值,第二个WebSocket客户端也能正常接收数据

测试用Python WebSocket客户端代码

import asyncio 
import random
import time
from websockets import connect
import json

async def func(uri):
  async with connect(uri) as websocket:
    await websocket.send("START-SEND-EVENTS")
    #await websocket.send("START-SEND-MESSAGES")
    message = await websocket.recv()
    print(message)
    while(True):
        
        msg = {
            "topic": "org.eclipse.ditto/camera01/things/twin/commands/modify",
            "headers": {
                "content-type": "text/plain"
            },
            "path": "features/coordinates/properties",
            "value": {"x": random.randrange(0,1000), "y": random.randrange(0,1000), "z": random.randrange(0,1000), "x_rotation": 0.0, "y_rotation": 0.0, "z_rotation": 0.0, "w_rotation": 1.0, "thingId": "org.eclipse.ditto:camera01"}
        }
        to_send = json.dumps(msg)
        time.sleep(1)
        await websocket.send(to_send)
        msg_recv = await websocket.recv()
        print(msg_recv)
uri = "ws://ditto:ditto@localhost:8080/ws/2"
asyncio.run(func(uri))

现有Kafka连接配置

{
"targetActorSelection": "/system/sharding/connection",
"headers": {
    "aggregate": false
},
"piggybackCommand": {
    "type": "connectivity.commands:modifyConnection",
    "connection": {
        "id": "kafka-connection-target",
        "connectionType": "kafka",
        "connectionStatus": "open",
        "failoverEnabled": true,
        "uri": "tcp://localhost:9092",
        "specificConfig":{
       "bootstrapServers":"localhost:9092"
         },
        "targets": [{
            "address": "topic_ditto",
            "topics": [
                "_/_/things/twin/events",
                "_/_/things/live/messages"
            ],
            "authorizationContext": ["ditto:unity"],
            "qos": 0
        }],
        "mappingContext": {
            "mappingEngine": "JavaScript",
            "options": {
                "incomingScript": "function mapToDittoProtocolMsg(headers, textPayload, bytePayload, contentType) {return null;}",
                "outgoingScript": "function mapFromDittoProtocolMsg(namespace, id, group, channel, criterion, action, path, dittoHeaders, value, status, extra) {let textPayload = '{\"x\":' + value.coordinates.properties.x + ',\"y\":' + value.coordinates.properties.y + ',\"z\":' + value.coordinates.properties.z + ',\"x_rotation\":' + value.coordinates.properties.x_rotation + ',\"y_rotation\": ' + value.coordinates.properties.y_rotation + ', \"z_rotation\": ' + value.coordinates.properties.z_rotation + ',\"w_rotation\":' + value.coordinates.properties.w_rotation + ',\"idCamera\":\"' + id + '\"}'; let bytePayload = null; let contentType = 'text/plain; charset=UTF-8'; return  Ditto.buildExternalMsg(dittoHeaders, textPayload, bytePayload, contentType);}",
                "loadBytebufferJS": "false",
                 "loadLongJS": "false"
            }
        }
    }
}
}

根因分析

问题由两个核心原因导致,共同造成WebSocket触发的更新无法投递到Kafka:

  1. 出站映射脚本逻辑错误
    当前出站JS脚本固定从value.coordinates.properties路径下取坐标字段,但Ditto事件中value的结构完全由修改的路径决定:
  • WebSocket发送的修改命令path是features/coordinates/properties,对应生成的孪生修改事件里,value直接就是坐标属性对象(即{x:xxx, y:xxx,...}结构),不存在coordinates这一层,脚本访问value.coordinates.properties时会直接抛出Cannot read properties of undefined的JS异常,这条消息会被Ditto直接丢弃,不会投递到Kafka。
  • Kafka源连接触发更新时,大概率修改的是Thing根节点或者features层级的路径,事件里的value包含完整的coordinates.properties层级结构,脚本能正常执行,所以消息可以正常投递到Kafka主题。
  1. 权限配置缺失(潜在问题)
    Kafka连接授权上下文用的是ditto:unity主体,默认通过WebSocket用ditto:ditto账号创建/修改的Thing,不会自动给ditto:unity授予READ权限,就算脚本修复后,Ditto也会因为权限校验不通过,过滤掉对应事件不会投递给Kafka连接。之前Kafka源连接更新时能正常投递,是因为源连接本身用ditto:unity身份操作,默认有权限读取自己触发的事件。

另外Python客户端代码存在不规范写法:在asyncio协程中使用同步阻塞的time.sleep(1),会阻塞整个事件循环,可能导致WebSocket消息收发异常;发送JSON内容时content-type错写为text/plain,不符合协议规范。

修复方案

  1. 修正Kafka目标配置与出站映射脚本
    首先在Kafka连接的target配置中添加extraFields配置,指定事件要携带的完整属性字段,避免修改路径变化导致value结构变动:
"targets": [{
    "address": "topic_ditto",
    "topics": [
        "_/_/things/twin/events",
        "_/_/things/live/messages"
    ],
    "authorizationContext": ["ditto:unity"],
    "qos": 0,
    "extraFields": ["features/coordinates/properties"]
}]

然后替换出站映射脚本,从extra携带的完整Thing快照中取坐标属性,用JSON.stringify替代手动字符串拼接避免格式错误,修改后的outgoingScript如下:

function mapFromDittoProtocolMsg(namespace, id, group, channel, criterion, action, path, dittoHeaders, value, status, extra) {
  const props = extra.features.coordinates.properties;
  const textPayload = JSON.stringify({
    "x": props.x,
    "y": props.y,
    "z": props.z,
    "x_rotation": props.x_rotation,
    "y_rotation": props.y_rotation,
    "z_rotation": props.z_rotation,
    "w_rotation": props.w_rotation,
    "idCamera": id
  });
  const bytePayload = null;
  const contentType = 'application/json; charset=UTF-8';
  return Ditto.buildExternalMsg(dittoHeaders, textPayload, bytePayload, contentType);
}
  1. 补全Thing权限配置
    给org.eclipse.ditto:camera01这个Thing添加策略,授予ditto:unity主体对Thing的READ权限,允许Kafka连接读取对应事件。
  2. 修正Python客户端代码
  • 把同步的time.sleep(1)替换为异步非阻塞的await asyncio.sleep(1)
  • 把请求头中的content-type改为application/json
    修正后的核心代码段:
msg = {
    "topic": "org.eclipse.ditto/camera01/things/twin/commands/modify",
    "headers": {
        "content-type": "application/json"
    },
    "path": "features/coordinates/properties",
    "value": {"x": random.randrange(0,1000), "y": random.randrange(0,1000), "z": random.randrange(0,1000), "x_rotation": 0.0, "y_rotation": 0.0, "z_rotation": 0.0, "w_rotation": 1.0, "thingId": "org.eclipse.ditto:camera01"}
}
to_send = json.dumps(msg)
await asyncio.sleep(1)
await websocket.send(to_send)

验证方式

配置修改后重启Kafka连接,再通过Python客户端发送更新指令,即可在Kafka的topic_ditto主题中正常接收到坐标更新消息。如果仍有异常,可以直接查看Ditto连接服务的日志,搜索kafka-connection-target相关的错误日志,会明确提示是脚本执行错误还是权限拒绝问题。


内容的提问来源于stack exchange,提问作者Lurù

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 14:48:34