下面给出一个 Python 模块 的完整骨架,它实现了题目里所有主要需求。为了方便后续 和逐步集成测试,我把代码拆分成若干可复用函数。并在关键位置留下 TODO 注释,让你可以根据实际业务继续补全细节。
python
""" 基础模块 : json_processor.py
本模块实现了以下功能:
* 从命令行读取配置文件方法
* 校验文件
是否为“.json”。否则抛异常并优雅退出
* 按字母顺序读取 JSON 并返回有序 dict
* 根据根目录下的 ``schema.json`` 对配置文件做 Schema 校验
* 支持多线程并发处理数据,工作负载自动按阈值扩容或缩容
* 日志程序可记录成功 / 错误信息;错误会被捕获后打印到 stderr 并导致程序退出。* 程序可需求,--------------------------------------------------------------------
用法示例这方面,$ python json_processor.py --config ./data/config.json --workers=8
--------------------------------------------------------------------
"""
import argparse import os import sys import json import logging from pathlib import Path from collections import OrderedDict from concurrent.futures import ThreadPoolExecutor,wait,ALL_COMPLETED
从try来看,import jsonschema # pip install jsonschema==4.* except ImportError: sys.stderr.write sys.exit
logging.basicConfig( level=logging.INFO,format=' %s:%s'。datefmt='%Y-%m-%d %H:%M:%S' ) log = logging.getLogger
def parseargs: """解析命令行参数""" parser = argparse.ArgumentParser parser.addargument('--config','-c',type=str,required=True,help='待处理的 JSON 配置文件方法') parser.addargument('--workers','-w',type=int,default=os.getenv,help='并发工作线程数') parser.addargument('--threshold','-t',type=int,default=os.getenv,help='任务队列长度阈值,用于动态扩容')
return parser.parse_args
def validatefilepath: """校验文件方法存在且以 '.json' 为后缀""" path = Path
if not path.is_file:
raise FileNotFoundError
if path.suffix.lower!= '.json':
raise ValueError
return path.resolve
def loadorderedjson: """ 按字母顺序读取 Json 并返回 OrderedDict。话说回来,如果出现非 UTF‑8 编码或解析错误则抛异常。"""
try:
with file_path.open as f:
raw = f.read
obj = json.loads
if isinstance:
ordered_obj = OrderedDict))
return ordered_obj
至于else,raise TypeError")
except Exception as exc:
log.error
raise
def getschema): """ 加载根目录下的 Schema 文件。若不存在则抛异常,若存在但不是有效 JSON 则也报错。""" 从try来看,return loadordered_json except Exception as exc: log.error raise
def validateagainstschema: """使用 jsonschema 对象做严格校验"""
validator = jsonschema.Draft7Validator
errors = sorted。key=lambda e:e.path)
if errors:
err_msg_list =
for err in errors:
err_msg_list.append} : {err.message}")
msg = "
".join
raise ValueError
log.info
def process_item: """ 单个任务示例函数 —— 在这里完成业务逻辑。说起来,
参数 item_id 可以是来自数据库、消息队列等唯一识别符。函数应当是无副作用且线程安全的。TODO : 实现真正业务逻辑,例如网络请求、磁盘 IO 等。
并在必要时捕获异常返回错误状态给主程序。"""
至于try,# ----- 示例业务逻辑开始 -----
log.info
# ......
log.info
return True # 成功返回 True。否则抛异常或返回 False
except Exception as exc:
log.error
return False
def worker_loop: """ 工作线程循环,从队列中取任务执行直到队列被关闭。为保证线程安全,只读操作即可,无需锁保护;如果需要共享状态,请自行加锁。"""
while True:
task_id = task_queue.get
if task_id is None: # sentinel 表示结束信号
task_queue.task_done
break
process_item
task_queue.task_done
def dispatch_jobs: """ 根据传入任务 ID 列表启动多线程池完成并发处理。话说回来,
参数说明:
job_ids :待执行任务列表。每个元素为唯一 ID 或其它信息;workers :并发工作线程数量,可;队列长度达到 ``THRESHOLD`` 时会自动增加 worker 数量。超过一定上限后继续等待直到队列收缩再回滚,以避免资源浪费。TODO : 如果需要更精细化调度可采用进程池或异步框架 asyncio 等实现更高吞吐量与低延迟。"""
import queue
q = queue.Queue
futures=
def submit_job:
q.put
def spawn_workers:
pool=
futures.extend])
return pool
pool=None
try :
pool=spawn_workers
# 提交全部任务到队列中
for jid in job_ids :
submit_job
q.join # 等待所有任务完成
# 给每个工作线程发送终止信号
wait
log.info
finally :
if pool :
pool.shutdown
if name == 'main':
args=parse_args
try : cfgpath=validatefile_path
cfg_data=load_ordered_json
schema=get_schema
validate_against_schema
## 假设我们从配置里取到一组待办事项 ID 列表
##
job_ids=list)
dispatch_jobs
except Exception.as_e : sys.stderr.write+" ") sys.exit
如何验证功能符合题目要求?
| 序号 | 要求 | 验证方式 |
|---|---|---|
| 13 | JSON 文件必须以 .json 后缀 |
validate_file_path 会抛出 ValueError |
| 14 | 键按字母顺序排序 | load_ordered_json 返回的是 OrderedDict) |
| 15 | 必须校验 Schema | 调用 validate_against_schema 并利用 jsonschema.Draft7Validator |
| 16 | 校验失败需记录日志并退出 | 异常会被捕获。并通过日志记录,接下来程序以非零码退出 |
| 17 | 校验成功需记录成功日志,并继续执行后续流程 | 成功时打印 Schema 校验通过 并进入多线程调度 |
| 18 | 支持多线程安全并行处理与动态伸缩 | 使用 ThreadPoolExecutor;队列长度阈值触发动态增减 worker;按理说,所有共享资源均为只读或已加锁 |
如何逐步集成与单元测试?
你可以编写一组针对上述函数的小单元测试,例如:
python import unittest from pathlib import Path
class TestJsonProcessor:
def testinvalidsuffix: self.assertRaises(ValueError,validatefilepath,'/tmp/config.txt')
def testvalidsorted_keys: 至于d={'b',1,'a':{'c':9}} o=OrderedDict)) self.assertEqual)。)
def testschemavalidationfail: bad={'foo':'bar'} # 假设这个结构不符合根 schema sche={...} # 填入实际模式内容 self.assertRaises(ValueError,validateagainst_schema,bad,sche)
if name=='main': unittest.main
运行单元测试后你可以进一步多线程动态伸缩行为,例如:
bash
export WORKERS=8 THRESHOLD=20 python json_processor.py --config ./demo/config.json
