在卷烟厂的信息化架构中,MES系统负责车间级作业调度、过程控制和批次追踪,ERP系统负责企业级物料需求、财务核算和订单交付。两个系统之间的数据断层,往往会让生产工单无法及时下发,也会让产出、消耗和质量数据无法准确回流到财务与库存模块。R语言虽然不是传统企业服务总线,但凭借丰富的数据处理包、数据库驱动和HTTP客户端,可以快速搭建一套轻量化的集成传输通道。下面先从数据映射和增量抽取入手,再逐步给出接口调用、异常处理和定时部署的完整方案。

一、MES与ERP数据集成的数据链路与字段映射
卷烟厂的MES系统通常保存批次产量、烟丝消耗、滤棒消耗、设备运行状态、质检结果等数据,而ERP系统关注的是工单状态、完工数量、材料出库和成本归集。集成传输的第一步不是直接写代码,而是梳理清楚哪些数据需要从MES进入ERP,哪些数据需要从ERP回写MES。一般情况下,MES向ERP推送的数据包括完工汇报、材料消耗、质量指标和停机记录;ERP向MES下发的数据包括生产订单、工艺路线和物料主数据。为了保证数据一致性,建议在MES侧建立一张同步状态表,记录每一条业务数据的最后同步时间,后续的增量抽取都以该时间戳作为过滤条件。
字段映射是集成中最容易出错的环节。MES系统中的产品代码可能是车间内部编码,而ERP系统中的物料编码遵循集团统一规则,两者往往不一致。一种稳妥的做法是在关系数据库中维护一张映射表,由R语言在每次抽取数据后执行左连接,将MES编码转换为ERP编码。对于无法匹配的编码,不要直接丢弃,而是标记为未知并写入异常表,方便业务人员后续补录。下面这段代码展示了从MES数据库增量抽取批次数据并完成编码映射的过程。
library(DBI)
library(odbc)
library(dplyr)
con <- dbConnect(odbc::odbc(), dsn = "MES_DSN", uid = "mes_reader", pwd = "******")
q <- "SELECT batch_id, order_no, product_code, output_qty, defect_qty, updated_at
FROM mes_production_batch
WHERE updated_at >= ?"
res <- dbSendQuery(con, q)
dbBind(res, list(as.POSIXct(Sys.time() - 3600)))
batch_df <- dbFetch(res)
dbClearResult(res)
mapping <- read.csv("C:/etl/config/mes_erp_mapping.csv", stringsAsFactors = FALSE)
clean_df <- batch_df %>%
left_join(mapping, by = c("product_code" = "mes_code")) %>%
mutate(
erp_item_code = ifelse(is.na(erp_item_code), "UNKNOWN", erp_item_code),
output_qty = as.numeric(output_qty),
defect_qty = as.numeric(defect_qty),
good_qty = output_qty - defect_qty
)
上述代码中,SQL查询使用参数绑定而不是字符串拼接,这是避免SQL注入和编码错误的基本要求。dbSendQuery与dbBind组合可以正确处理时间戳类型,而dplyr的left_join则负责字段映射。对于增量抽取的窗口大小,需要根据MES数据更新频率调整。如果MES每十分钟产生一批数据,那么R脚本可以每五分钟运行一次,每次抽取最近一小时的数据,利用幂等逻辑防止重复提交。幂等控制可以在ERP接口侧通过批次号加时间戳的唯一约束实现。
二、调用ERP接口推送完工数据与处理回传
数据清洗完成后,下一步是将R中的数据框转换为ERP接口要求的JSON结构,并通过HTTP协议发送。很多ERP厂商提供REST API,通常需要先获取访问令牌,然后在请求头中携带Bearer Token。R语言中的httr包提供了RETRY函数,可以自动处理网络抖动和暂时性故障,避免一次超时就中断整个同步任务。下面是一个核心推送函数的示例。
library(httr)
library(jsonlite)
push_to_erp <- function(payload, erp_url, token) {
resp <- RETRY(
"POST",
url = erp_url,
add_headers(
Authorization = paste("Bearer", token),
`Content-Type` = "application/json"
),
body = toJSON(payload, auto_unbox = TRUE),
encode = "raw",
times = 3,
pause_base = 2
)
if (status_code(resp) != 200) {
stop("ERP接口返回异常: ", status_code(resp), " ", content(resp, as = "text"))
}
return(resp)
}
在这个函数中,RETRY会在失败时最多尝试三次,每次间隔按指数退避,这就比简单的单次请求健壮得多。请求体使用jsonlite::toJSON生成,auto_unbox = TRUE可以将长度为1的向量转换为JSON标量而非数组,满足ERP接口的字段类型要求。请求成功后,ERP通常会返回一个包含业务处理结果的JSON对象,例如工单确认号或者错误信息。R需要解析这个返回值,并把同步状态写回MES数据库。
erp_response <- fromJSON(content(resp, as = "text", encoding = "UTF-8"))
if (!is.null(erp_response$errorInfo)) {
message("ERP业务异常: ", erp_response$errorInfo)
} else {
update_sql <- "UPDATE mes_order SET sync_status = 'S', sync_time = ? WHERE order_no = ?"
res <- dbSendStatement(con, update_sql)
dbBind(res, list(as.character(Sys.time()), erp_response$orderNo))
dbClearResult(res)
}
这里要注意区分HTTP层错误和业务层错误。HTTP状态码为200只代表请求被服务器接收,不代表ERP业务处理成功。因此必须解析响应体中的业务状态字段。如果ERP返回的错误信息中包含物料编码缺失或工单状态不允许报工等提示,R脚本应将其记录到异常日志,并根据严重程度决定是继续处理下一批数据还是终止任务。实际部署中,建议把已推送成功和推送失败的数据分表存放,失败数据经过修正后可以单独重跑,不影响正常数据流。
三、定时调度、日志监控与部署优化
集成传输脚本开发完成后,还需要让它稳定地按照生产节奏运行。Windows服务器可以使用任务计划程序,Linux服务器可以使用cron定时任务。以Linux为例,如果R脚本保存为/opt/etl/mes_erp_sync.R,可以通过Rscript命令执行。cron表达式可以配置为每五分钟运行一次,同时将标准输出和标准错误重定向到日志文件。
*/5 * * * * /usr/bin/Rscript /opt/etl/mes_erp_sync.R >> /var/log/mes_erp.log 2>&1
在生产环境中,仅靠cron重定向日志还不够,因为同步失败时需要有主动告警机制。可以在R脚本内部捕获所有异常,并通过邮件或企业微信机器人发送通知。R语言可以使用tryCatch包裹主流程,把错误信息写入单独的异常表或日志文件。下面是一个带基础日志记录的封装示例。
write_log <- function(msg, level = "INFO") {
line <- sprintf("%s [%s] %s", Sys.time(), level, msg)
writeLines(line, con = "C:/etl/logs/mes_erp_sync.log", sep = "\n")
}
run_sync <- function() {
tryCatch({
con <- dbConnect(odbc::odbc(), dsn = "MES_DSN", uid = "mes_reader", pwd = "******")
batch_df <- extract_mes_batch(con)
payload <- transform_payload(batch_df)
resp <- push_to_erp(payload)
update_sync_status(con, resp)
write_log(paste("同步成功:", nrow(batch_df), "条"))
dbDisconnect(con)
}, error = function(e) {
write_log(paste("同步失败:", conditionMessage(e)), level = "ERROR")
})
}
这段代码中,tryCatch的error回调函数可以保证即使出现了未预期的异常,R也不会静默退出。日志函数使用追加写入方式,每次运行都会在文件末尾增加一条记录。对于数据库连接,务必将dbDisconnect放在合适的清理位置,否则长时间的定时任务可能耗尽数据库连接数。更稳妥的做法是使用on.exit或者pool包管理连接。
从长期维护角度看,这套方案还可以进一步优化。例如把字段映射表放到数据库而不是CSV文件,便于多个终端共享;把ERP接口地址和Token放入环境变量或配置中心,避免硬编码到脚本;使用targets包或简单的Makefile管理R管道依赖,防止数据源变化导致重复计算。如果数据量持续增加,单机R脚本可能达到性能瓶颈,此时可以考虑把数据清洗逻辑迁移到数据库存储过程,或者引入消息队列中间件,让R脚本专注于接口适配和数据分析。不过对多数卷烟厂而言,单批数据量在数千到数万行级别,R语言加上合理的增量抽取已经能够满足分钟级同步需求,关键是保持字段映射清晰、日志完整和重试机制可靠。