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:
- 出站映射脚本逻辑错误
当前出站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主题。
- 权限配置缺失(潜在问题)
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,不符合协议规范。
修复方案
- 修正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); }
- 补全Thing权限配置
给org.eclipse.ditto:camera01这个Thing添加策略,授予ditto:unity主体对Thing的READ权限,允许Kafka连接读取对应事件。 - 修正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ù
相关产品推荐
相关产品推荐

