
这次我们来看一个在 DolphinDB 中处理批处理作业的实用功能。对于需要定时、批量执行数据清洗、报表生成或模型训练等任务的数据工程师和开发者来说这是一个能极大提升自动化效率的核心特性。它允许你将复杂的脚本封装成可调度的作业在指定时间或周期性地自动运行从而解放人力确保数据处理流程的稳定与准时。本文将聚焦于 DolphinDB 批处理作业的基础概念与核心操作。我们会先快速了解它的核心能力与适用场景然后通过具体的环境准备、脚本编写、作业提交与管理的全流程带你完成一次完整的批处理作业实践。无论你是想自动化每日的 KPI 计算还是需要定时执行高频因子计算这篇文章都能提供清晰的指引。1. 核心能力速览在深入细节之前我们先通过一个表格快速把握 DolphinDB 批处理作业的核心特性这有助于你判断它是否适合解决你手头的问题。能力项说明项目类型数据库内置的作业调度与执行引擎。核心功能提交、调度并执行后台批处理脚本支持定时一次性/周期性和立即执行。执行环境在 DolphinDB 服务端运行与数据库引擎深度集成可直接操作库表数据。资源管理作业在后台运行不阻塞前端会话可通过系统视图监控作业状态与资源占用。依赖管理支持通过submitJob命名参数传递上下文变量作业脚本可访问提交时的局部变量。适合场景定时数据导入/清洗、日终报表计算、定期模型再训练、流计算任务启停等自动化流程。2. 适用场景与使用边界DolphinDB 的批处理作业功能并非一个独立的外部调度系统而是其数据库内核提供的一种将脚本“后台化”和“计划化”执行的能力。理解其适用场景和边界能帮助你更有效地利用它。它非常适合以下场景定时ETL任务每天凌晨定时从外部数据源如文件、其他数据库导入数据并进行清洗、转换后存入 DolphinDB。定期报表生成在交易结束后自动计算当日各类业务指标生成报表并导出或写入结果表。模型批量训练与预测对于机器学习或量化因子模型定期如每周使用新增数据重新训练模型并批量生成预测结果。维护任务定期清理历史数据、备份数据库、或更新分布式数据库的维度表。配合流计算在特定时间点自动启停流数据订阅、或修改流计算引擎的参数。需要注意的使用边界非独立调度系统它缺乏像 Airflow 那样复杂的依赖调度、任务编排 DAG 可视化界面。复杂的跨作业依赖需要自行通过状态表或文件来协调。资源隔离有限作业共享数据库服务器的计算资源CPU、内存。大量并发或资源密集型作业可能相互影响需要根据服务器能力合理规划。执行环境作业脚本在服务端执行其权限与提交作业的用户权限一致。脚本中无法直接访问客户端本地文件系统除非通过函数如loadText读取服务端可访问路径的文件。生命周期作业由 DolphinDB 服务管理。服务重启后未完成的作业会终止定时作业需要重新提交除非使用持久化的定时任务功能。3. 环境准备与前置条件要开始使用批处理作业你需要一个正在运行的 DolphinDB 环境。以下是通用的环境检查清单DolphinDB 服务确保 DolphinDB 单机版或集群版已安装并启动。这是作业运行的基础平台。客户端连接工具你需要通过以下任一方式连接至 DolphinDB 服务DolphinDB GUI官方图形化客户端适合交互式开发和作业管理。VS Code 插件使用 DolphinDB 插件进行脚本开发。Python/Java/C API通过编程接口连接适合集成到自动化系统中提交作业。命令行工具通过dolphindb命令行连接。用户权限确保你使用的登录账号拥有执行目标操作如读写特定数据库、表执行函数的权限。作业将以提交用户的身份运行。脚本知识熟悉 DolphinDB 脚本语言的基本语法包括变量定义、函数编写、SQL 查询以及控制流语句。4. 核心函数与作业提交DolphinDB 提供了几个核心函数来管理批处理作业。我们重点看最常用的submitJob函数。4.1submitJob函数详解submitJob用于提交一个后台作业并立即返回一个作业 ID而不等待作业执行完成。submitJob(jobId, jobDesc, function, args...)jobId一个字符串用于标识作业。如果已存在同 ID 的作业提交会失败。jobDesc作业描述便于后续识别。function要在后台执行的函数名或函数定义。args...可变参数传递给function的参数。关键特性通过submitJob提交时函数及其参数会被序列化并传到服务端。函数内部可以访问提交时刻的局部变量通过命名参数捕获但无法感知作业提交后客户端发生的变量变化。4.2 一个简单的作业提交示例假设我们有一个简单的需求将某个表的数据按日聚合后保存到新表。我们可以这样封装和提交作业// 1. 定义要在后台执行的函数 def dailyAggregation(tableName, outputTableName){ // 获取昨天的日期 targetDate today() - 1 // 执行聚合查询这里假设表中有tradeTime和price字段 sqlText select date(tradeTime) as tradeDate, sum(price) as dailySum from ${tableName} where date(tradeTime)${targetDate} group by date(tradeTime) result sql(sqlText) // 将结果插入到输出表假设表已存在 outputTableName.append!(result) return Finished aggregation for targetDate } // 2. 准备参数这些局部变量会被捕获 inputTable trades outputTable daily_summary // 3. 提交批处理作业 jobId agg_ string(today()) jobDesc Daily trade sum aggregation jobId submitJob(jobId, jobDesc, dailyAggregation, inputTable, outputTable) // 打印作业ID用于后续查询 print(Submitted job with ID: jobId)执行上述代码后dailyAggregation函数会在数据库服务后台异步执行你的当前会话可以继续处理其他任务。5. 作业状态监控与管理作业提交后如何知道它是否成功、正在运行还是已经失败DolphinDB 提供了系统视图供你查询。5.1 查询作业状态最常用的系统视图是getRecentJobs()和getJobStatus。getRecentJobs()返回近期所有作业的列表包括已完成和正在运行的。// 查看最近的所有作业 jobs getRecentJobs() select * from jobsgetJobStatus(jobId)获取指定作业 ID 的详细状态。// 查询我们刚才提交的作业状态 status getJobStatus(jobId) print(status)getJobStatus返回一个字典包含如jobId、jobDesc、startTime、endTime、errorMsg、priority、parallelism等关键信息。其中errorMsg在作业失败时会包含错误信息是排查问题的首要位置。5.2 作业管理操作除了查询你还可以对作业进行一些控制取消作业如果作业挂起或你想中止一个长时间运行的任务可以使用cancelJob(jobId)。cancelJob(jobId)获取作业返回结果如果作业函数有返回值可以通过getJobReturn(jobId, [blockingfalse])获取。blocking参数为true时会等待作业完成。// 非阻塞式获取如果作业未完成则抛出异常 try{ result getJobReturn(jobId) } catch(ex){ print(Job not finished yet.) } // 阻塞式获取等待作业完成 result getJobReturn(jobId, true)6. 定时作业scheduleJob的使用对于需要周期性执行的任务使用submitJob每次手动提交并不现实。这时就需要scheduleJob函数。6.1 创建定时作业scheduleJob可以创建一次性或周期性的定时作业。scheduleJob(jobId, jobDesc, scheduleTime, function, args...)scheduleTime一个时间标量或向量指定作业的执行时间。可以是绝对时间如2024.05.01 00:00:00也可以是相对时间字符串如00:00:00表示每天零点。其他参数与submitJob类似。示例1创建一次性的定时作业// 在明天上午10点执行一次数据备份 scheduleTime now() 1*24*60*60*1000 // 当前时间加1天毫秒 scheduleJob(backup_once, One-time backup, scheduleTime, backupFunction)示例2创建周期性的定时作业每日执行// 每天凌晨2点执行数据清理任务 scheduleTime 02:00:00 scheduleJob(daily_clean, Daily data cleanup, scheduleTime, cleanupFunction, log_table)6.2 管理定时作业查询所有定时作业使用getScheduledJobs()。select * from getScheduledJobs()删除定时作业使用deleteScheduledJob(jobId)。deleteScheduledJob(daily_clean)注意删除的是未来的调度计划不会影响已经提交到执行队列中的作业实例。7. 资源占用与性能观察批处理作业在后台运行其资源消耗直接影响数据库服务的整体性能。你需要知道如何观察和评估。通过系统函数观察getJobStatus查看作业的startTime和endTime可以计算其运行时长。getRecentJobs观察作业的priority和parallelism设置如果使用了这些特性。通过操作系统工具观察由于作业线程是 DolphinDB 进程的一部分你可以通过top(Linux)、任务管理器(Windows) 或htop等工具监控 DolphinDB 进程的 CPU 和内存使用率。一个长时间高占用的作业会显著提升进程的资源使用。性能考量I/O 密集型作业如大数据量导入导出会显著增加磁盘 I/O可能影响同时进行的查询响应速度。CPU 密集型作业如复杂的因子计算、模型训练会占用大量 CPU 时间片。内存密集型作业如对超大表进行全表扫描或分组聚合可能引发内存压力。建议将大型批处理作业安排在业务低峰期如夜间执行。对于集群环境可以考虑利用分布式计算能力将任务分发到多个数据节点并行执行。8. 常见问题与排查方法在实际使用中你可能会遇到以下问题。这里提供一个排查指南。问题现象可能原因排查方式解决方案作业提交失败提示jobId already exists重复提交了相同jobId的作业。检查jobId是否唯一。使用更唯一的标识符如结合时间戳jobId myJob_ string(now())。作业状态一直是running长时间不结束。1. 作业逻辑有死循环。2. 处理的数据量极大确实需要长时间运行。3. 作业在等待某个锁如写表锁。1. 检查作业函数逻辑。2. 通过getJobStatus查看startTime。3. 检查相关表是否有未结束的事务。1. 优化脚本逻辑或增加中断条件。2. 如果正常则等待。3. 排查并结束阻塞的事务。使用cancelJob强制中止异常作业。作业失败errorMsg显示权限错误。提交作业的用户对作业中访问的数据库、表或函数缺乏足够权限。仔细阅读errorMsg确认是哪个操作被拒绝。使用更高权限的用户提交作业或联系管理员为当前用户授予相应权限。定时作业没有按预期时间执行。1.scheduleTime格式错误或已过时。2. DolphinDB 服务在计划时间点重启了。3. 定时作业被意外删除。1. 使用getScheduledJobs()确认作业计划是否存在及下次执行时间。2. 检查服务日志看是否有重启记录。1. 确保scheduleTime是未来的绝对时间或合法的相对时间字符串。2. 对于关键任务考虑在服务启动脚本中重新提交定时作业。作业函数中访问的变量值为NULL或未定义。作业函数只能捕获提交时的局部变量。如果变量是后续赋值或来自全局作用域可能无法正确传递。检查submitJob或scheduleJob调用处的变量作用域。将所有依赖的参数都通过args...显式传递给作业函数避免依赖外部变量。通过 API如 Python提交作业后客户端断开连接作业是否继续作业在服务端执行与客户端连接无关。在 DolphinDB GUI 或另一个会话中查询getRecentJobs()。作业会继续执行直至完成或出错。这是批处理作业的核心优势之一。9. 最佳实践与使用建议为了更稳健、高效地使用批处理作业遵循以下实践会大有裨益作业脚本模块化与测试将计划在作业中执行的逻辑封装成独立的函数并在前端会话中充分测试其正确性和性能。确保函数内部有完善的异常处理try-catch并将错误信息记录到日志表或文件中而不是让作业静默失败。作业标识与日志为jobId和jobDesc设计清晰的命名规则例如{功能}_{日期}_{序列号}。在作业函数内部将关键步骤、开始结束时间、处理行数、错误信息等写入一个专用的作业日志表。这比单纯依赖系统视图getRecentJobs更利于长期追踪和审计。资源与调度规划评估作业的资源消耗避免在业务高峰时段运行重型作业。对于多个有关联的作业不要依赖隐式的执行顺序。可以通过在第一个作业完成后向一个状态表写入标记第二个作业轮询这个标记的方式来实现简单依赖。参数化与配置化避免将数据库名、表名、文件路径等硬编码在作业函数中。可以考虑通过一个配置表或参数文件来管理这些变量作业函数运行时去读取。安全与权限遵循最小权限原则使用仅具备作业所需操作权限的专用账号来提交作业。定期审查定时作业列表 (getScheduledJobs())清理不再需要的作业。掌握 DolphinDB 的批处理作业功能相当于为你数据处理的流水线装上了自动化的齿轮。从一次性的数据迁移到日复一日的指标计算它都能可靠地在后台完成任务。建议你从一个小而具体的任务开始尝试比如定时计算某个数据表的行数并记录逐步熟悉整个“定义函数 - 提交作业 - 监控状态 - 查看结果”的流程。当你成功跑通第一个自动化作业后自然会想到将其应用到更多重复性的工作中去。