ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

WMS与WCS任务下发代码包:出库入库接口对接与重试对账实战

WMS与WCS任务下发代码包:出库入库接口对接与重试对账实战 简介本资源为WMS与WCS系统对接的通信代码示例聚焦仓储管理系统中WMS向WCS下发任务的JSON报文格式与字段定义面向自动化立体库、智能仓储方向的开发人员与集成工程师。资源以cmd指令区分入库、出库、移库三类业务并完整给出seq序号、task_id任务唯一编码、起止站台、起止排/列/层坐标及重量、条码等字段的注释说明便于读者理解任务调度接口的数据结构。压缩包共约2000个文件以dll动态库、png图片、cshtml页面、js脚本、cs与css样式文件为主另含xml、config配置及pdb调试文件整体约367MB属于典型.NET Web项目结构。目前已有2355人学习下载可作为WCS任务下发模块的参考实现帮助读者快速掌握报文组装、字段含义与接口联调思路。1. 从 WMS 到 WCS 的任务下发一份能跑通的代码包到底长什么样很多做仓储系统的兄弟第一次接触 WCS 对接都会卡在同一个地方WMS 里出库单、入库单、波次都跑通了可任务就是下不到设备侧堆垛机、AGV、输送线全在原地等。这份代码包解决的就是这个断点——它把 WMS 生成出库/入库任务后如何组装报文、如何调用 WCS 接口、如何处理回执与重试完整落成了一套可复现的源码。适合两类人一是刚接手 WMS-WCS 接口层的后端二是需要给现场设备联调提供稳定任务流的实施。它不讲仓储理论只讲任务从 WMS 数据库里被捞出来、变成 WCS 能认的指令、再被确认执行这一条链路。下面按我拆包的顺序把选型、代码、参数和坑一次说清。2. 任务模型与接口选型为什么不是直接写 WCS 库表2.1 WMS 与 WCS 的职责边界先把边界划清楚不然后面代码全是糊的。WMS 管的是「要做什么」——出库单、入库单、库存分配、波次、任务优先级WCS 管的是「怎么动」——把任务拆成设备可执行的动作序列调度堆垛机、穿梭车、输送线、AGV。两者之间传递的核心对象就是任务Task出库任务和入库任务是两条最常走的链路。常见做法是 WMS 生成任务后不直接写 WCS 的库表而是通过接口下发。原因有三个一是解耦WCS 换供应商时 WMS 不用重写二是可追溯每次下发都有报文和回执三是可控WCS 侧能对任务做校验、限流、排队。直接写库表看起来快但一旦 WCS 表结构变动或者需要加校验WMS 就被拖死这是血泪经验。这份代码包采用的模型是WMS 侧维护一张任务表任务状态从「待下发」到「已下发」到「执行中」到「完成/失败」WCS 侧提供任务接收接口和状态回传接口。出库任务和入库任务共用一套下发框架差异在任务类型字段和扩展参数上。2.2 接口协议与报文格式的取舍接口选型上常见的有 REST/JSON、WebService、MQ 消息、TCP 自定义报文。这份代码包用的是 REST/JSON 为主、MQ 为辅的结构任务下发走 HTTP POST状态回传走 MQ 或者回调接口。为什么这么选因为现场联调时 HTTP 最容易抓包和模拟用 Postman 就能打而状态回传用 MQ 能避免 WMS 被 WCS 的回调打爆。报文格式上出库任务和入库任务的字段有共性也有差异。共性字段包括任务号、任务类型、优先级、容器号、起始位置、目标位置、创建时间出库任务多一个出库单号和波次号入库任务多一个入库单号和供应商信息。代码包里把这些抽成一个基类子类只扩展差异字段避免每个任务类型写一遍。提示接口字段命名一定要和 WCS 供应商对齐尤其是位置编码格式。我见过因为 WMS 用「A-01-02-03」而 WCS 要「A010203」导致任务全部下不去的情况联调前先对字段字典。2.3 任务下发的整体流程流程拆成六步WMS 业务触发任务生成、任务落库为待下发、下发服务捞取待下发任务、组装报文调用 WCS 接口、根据回执更新任务状态、失败进入重试队列。这六步在代码里对应六个模块下面会逐个落到代码。3. 出库任务下发从任务组装到接口调用的完整代码3.1 任务实体与状态机设计先看任务实体。出库任务和入库任务共用一个 Task 基类状态用枚举管理避免到处写魔法数字。from enum import Enum from dataclasses import dataclass, field from datetime import datetime class TaskType(Enum): OUTBOUND OUTBOUND # 出库任务 INBOUND INBOUND # 入库任务 class TaskStatus(Enum): PENDING PENDING # 待下发 SENT SENT # 已下发等回执 EXECUTING EXECUTING # WCS 已接收执行中 DONE DONE # 完成 FAILED FAILED # 失败可重试 dataclass class Task: task_no: str # WMS 侧任务号全局唯一 task_type: TaskType priority: int 5 # 1 最高9 最低 container_no: str # 容器/托盘号 from_location: str # 起始位置编码 to_location: str # 目标位置编码 status: TaskStatus TaskStatus.PENDING retry_count: int 0 created_at: datetime field(default_factorydatetime.now) ext: dict field(default_factorydict) # 扩展字段放单号、波次等这段代码的关键点是task_no必须由 WMS 生成且全局唯一WCS 侧用它做幂等。ext字典用来放差异字段出库放outbound_no和wave_no入库放inbound_no和supplier。retry_count是重试次数的计数器后面重试逻辑会读它。状态机的流转规则要写死PENDING 只能到 SENTSENT 到 EXECUTING 或 FAILEDEXECUTING 到 DONE 或 FAILEDFAILED 可以回到 PENDING 重试。不要允许跳状态否则现场排查时根本不知道任务卡在哪。3.2 出库任务组装与下发代码出库任务的核心是把 WMS 的出库单信息转成 WCS 能认的指令。下面这段是组装和下发的主逻辑。import requests import json WCS_TASK_URL http://wcs-host:8080/api/task/receive TIMEOUT 5 # 秒现场网络抖动大别设太长 def build_outbound_payload(task: Task) - dict: 把出库任务组装成 WCS 报文 return { taskNo: task.task_no, taskType: OUT, priority: task.priority, containerNo: task.container_no, fromLocation: task.from_location, toLocation: task.to_location, outboundNo: task.ext.get(outbound_no, ), waveNo: task.ext.get(wave_no, ), createTime: task.created_at.strftime(%Y-%m-%d %H:%M:%S) } def send_task_to_wcs(task: Task) - bool: 下发单个任务返回是否成功 payload build_outbound_payload(task) try: resp requests.post( WCS_TASK_URL, datajson.dumps(payload), headers{Content-Type: application/json}, timeoutTIMEOUT ) if resp.status_code 200: result resp.json() # WCS 约定 code0 表示接收成功 if result.get(code) 0: task.status TaskStatus.SENT return True else: # 业务失败记录 WCS 返回的错误信息 task.ext[wcs_error] result.get(msg, unknown) return False return False except requests.exceptions.Timeout: task.ext[wcs_error] timeout return False except requests.exceptions.RequestException as e: task.ext[wcs_error] str(e) return False逻辑说明build_outbound_payload负责字段映射注意taskType用的是 WCS 约定的「OUT」而不是枚举值这种映射一定要单独抽函数别散在业务代码里。send_task_to_wcs里区分了 HTTP 层失败和业务层失败超时单独捕获因为超时和业务拒绝的处理策略不同——超时可能是任务已经进去了但回执没回来需要靠幂等查询确认不能直接重发。参数说明WCS_TASK_URL是 WCS 提供的任务接收地址现场部署时换成实际 IP 和端口。TIMEOUT设 5 秒是因为 WCS 接收任务通常是写库或入队正常在毫秒级超过 5 秒基本是网络或 WCS 卡了。priority直接透传WCS 侧按这个排序。3.3 批量下发与并发控制单个下发跑通后实际场景是批量。出库波次一放就是几十上百个任务不能串行一个个发但也不能无脑并发把 WCS 打挂。from concurrent.futures import ThreadPoolExecutor, as_completed MAX_WORKERS 8 # 并发数按 WCS 承受能力调 def batch_send(tasks: list) - dict: 批量下发返回成功和失败列表 success, failed [], [] with ThreadPoolExecutor(max_workersMAX_WORKERS) as executor: future_map {executor.submit(send_task_to_wcs, t): t for t in tasks} for future in as_completed(future_map): task future_map[future] if future.result(): success.append(task.task_no) else: failed.append(task.task_no) return {success: success, failed: failed}并发数MAX_WORKERS是现场调出来的默认 8。WCS 如果是单线程接收并发高了反而排队超时这时候要降到 2 到 4。批量下发后成功和失败要分别落库失败的进重试队列不要直接丢弃。4. 入库任务与状态回传回执处理、幂等与重试机制4.1 入库任务的差异处理入库任务和出库任务共用下发框架差异在报文组装。入库多的是入库单号、供应商、质检状态起始位置通常是收货口目标位置是存储位。def build_inbound_payload(task: Task) - dict: 入库任务报文注意 taskType 和扩展字段 return { taskNo: task.task_no, taskType: IN, priority: task.priority, containerNo: task.container_no, fromLocation: task.from_location, toLocation: task.to_location, inboundNo: task.ext.get(inbound_no, ), supplier: task.ext.get(supplier, ), qcStatus: task.ext.get(qc_status, PASS), createTime: task.created_at.strftime(%Y-%m-%d %H:%M:%S) }qcStatus是质检状态入库时如果质检未完成WCS 可能要把任务送到待检区而不是存储区这个字段一定要和 WCS 对齐取值。supplier在部分 WCS 里用于分区策略别漏传。4.2 状态回传接口与幂等处理WCS 执行完任务后会回传状态WMS 侧要提供一个接收接口。这个接口最大的坑是重复回传所以必须幂等。def handle_wcs_callback(data: dict) - dict: 处理 WCS 状态回传幂等 task_no data.get(taskNo) wcs_status data.get(status) # EXECUTING / DONE / FAILED if not task_no: return {code: 1, msg: taskNo missing} task load_task(task_no) if task is None: return {code: 1, msg: task not found} # 幂等已完成的任务不再处理 if task.status TaskStatus.DONE: return {code: 0, msg: already done} if wcs_status EXECUTING: task.status TaskStatus.EXECUTING elif wcs_status DONE: task.status TaskStatus.DONE elif wcs_status FAILED: task.status TaskStatus.FAILED task.ext[wcs_error] data.get(errorMsg, ) save_task(task) return {code: 0, msg: ok}幂等判断放在最前面DONE状态直接返回成功避免重复更新。load_task和save_task是持久层方法实际项目里换成 ORM 或 DAO。回传接口要记录原始报文现场扯皮时这是唯一证据。4.3 重试机制与死信处理失败任务不能无限重试要有上限和退避。MAX_RETRY 3 RETRY_INTERVAL [10, 30, 60] # 秒退避间隔 def retry_failed_tasks(): 捞取失败任务重试 tasks query_tasks(statusTaskStatus.FAILED) for task in tasks: if task.retry_count MAX_RETRY: mark_dead_letter(task) # 进死信人工处理 continue task.retry_count 1 task.status TaskStatus.PENDING save_task(task) # 实际下发由调度器按 RETRY_INTERVAL 触发MAX_RETRY设 3 是经验值超过 3 次基本是数据问题不是网络问题再重试也是浪费。退避间隔按retry_count取第一次 10 秒第二次 30 秒第三次 60 秒。死信任务要能人工干预现场经常需要手动改位置或容器号后重新下发。注意重试前一定要确认 WCS 侧没有已经接收该任务。如果第一次下发是超时但 WCS 实际收到了重试会导致重复任务。常见做法是重试前先调 WCS 的任务查询接口确认。5. 避坑与排查任务下发不成功时先看这几处5.1 任务一直停在 PENDING 不下发现象任务落库了状态是 PENDING但下发服务不捞。原因通常是调度器没启动或者捞取条件写错比如只捞了出库没捞入库或者时间条件把新任务过滤掉了。解决先看调度器日志有没有执行再检查捞取 SQL 的 where 条件把状态和时间范围打出来。5.2 WCS 返回成功但任务没执行现象下发接口返回 code0但设备不动任务状态停在 SENT。原因多半是 WCS 接收后入队了但调度没触发或者位置编码 WCS 不认识导致任务被挂起。解决让 WCS 侧查任务队列确认任务是否被消费同时核对位置编码格式这是最高频的翻车点。5.3 重复任务导致设备动作两次现象同一个容器被送了两次或者堆垛机执行了重复指令。原因是没有幂等超时重试或者 WCS 重复回传都会触发。解决WMS 侧下发前用 task_no 查重WCS 侧接收时也做幂等回传接口对 DONE 状态直接返回成功不重复处理。5.4 状态回传丢失导致任务永远 EXECUTING现象任务卡在 EXECUTINGWCS 说已完成WMS 没收到回传。原因是 MQ 丢消息或者回调接口报错。解决加对账机制定时用 WCS 的任务查询接口拉取状态和 WMS 本地比对不一致的以 WCS 为准修正。这是黑匣子最容易被忽略的地方。5.5 批量下发把 WCS 打挂现象波次一放WCS 接口超时任务大面积失败。原因是并发太高WCS 接收能力有限。解决降MAX_WORKERS加限流或者改成 MQ 异步下发让 WCS 自己消费。现场联调时先小批量试别一上来就放全量波次。6. 进阶用对账任务兜住状态不一致附验证脚本前面讲的都是正常链路但现场跑久了状态不一致是必然的。WMS 和 WCS 各有一套状态网络抖动、服务重启、MQ 丢消息都会让两边对不上。我一般会加一个对账任务定时跑把差异捞出来修正。这是这套代码包里最值钱的部分因为它决定了系统能不能长期稳定跑。对账的逻辑是拉取 WMS 侧所有非终态任务PENDING、SENT、EXECUTING逐个调 WCS 的任务查询接口拿到 WCS 侧状态后比对。WCS 说 DONE 而 WMS 还是 EXECUTING 的修正为 DONEWCS 说没有这个任务的说明下发根本没成功回退到 PENDING 重试WCS 说 FAILED 的同步失败原因。def reconcile_tasks(): 对账以 WCS 状态为准修正 WMS wms_tasks query_tasks(status_in[ TaskStatus.PENDING, TaskStatus.SENT, TaskStatus.EXECUTING ]) fixed 0 for task in wms_tasks: wcs_info query_wcs_task(task.task_no) # 调 WCS 查询接口 if wcs_info is None: # WCS 没有该任务回退重试 task.status TaskStatus.PENDING task.retry_count 1 save_task(task) fixed 1 continue wcs_status wcs_info.get(status) if wcs_status DONE and task.status ! TaskStatus.DONE: task.status TaskStatus.DONE save_task(task) fixed 1 elif wcs_status FAILED and task.status ! TaskStatus.FAILED: task.status TaskStatus.FAILED task.ext[wcs_error] wcs_info.get(errorMsg, ) save_task(task) fixed 1 return fixedquery_wcs_task是 WCS 提供的任务查询接口如果 WCS 没提供这个方案就落不了地所以接口选型时一定要把查询接口写进合同。对账频率我一般设 5 分钟一次太频繁浪费资源太稀疏差异积累多了修正成本高。对账结果要打日志修正了多少条、哪些任务被回退这些数据是判断系统健康度的指标。验证脚本可以单独跑不依赖调度器# 手动触发一次对账观察输出 python -m wms_wcs.reconcile --once --verbose # 输出示例reconcile done, fixed3, pending_retry1, done_sync2跑完看fixed数量如果每次都是 0说明链路健康如果持续大于 0说明回传链路有问题要去查 MQ 或回调接口。我习惯在每次现场上线后连续观察三天对账数据稳定为 0 才算验收通过。从那以后我每次做 WMS-WCS 对接都会先把对账脚本写好再写下发逻辑因为下发只是开始能长期对得上才是终点。希望帮到你。本文还有配套的精品资源点击获取
返回列表