Skip to main content

A library for processing JSON data with operators and pipelines

Project description

JSONFlow

一个高效、灵活的JSON数据流式处理库,专为大规模数据处理和大语言模型交互设计。

核心特性

  • 流式处理架构: 将JSON数据通过操作符链式处理,支持单个对象和批量处理
  • 集合处理能力: 智能处理JSON列表,支持自动展平或保持嵌套结构
  • 大语言模型集成: 专门的模型操作符,轻松调用各种大语言模型
  • 并发执行: 多线程/多进程并行处理,保持输入顺序
  • 字段透传: 自动保留指定字段,避免重复处理系统字段
  • 系统字段管理: 简化ID、时间戳等系统级字段的添加与管理
  • 丰富的操作符: 内置文本处理、字段操作、表达式计算等多种操作符
  • 操作符IO日志: 全面的日志系统,便于调试和开发

安装

pip install jsonflow

快速开始

基本用法

from jsonflow.core import Pipeline
from jsonflow.io import JsonLoader, JsonSaver
from jsonflow.operators.json_ops import TextNormalizer
from jsonflow.operators.model import ModelInvoker

# 创建处理管道
pipeline = Pipeline([
    TextNormalizer(),
    ModelInvoker(model="gpt-3.5-turbo"),
])

# 处理单个JSON文件
loader = JsonLoader("input.jsonl")
saver = JsonSaver("output.jsonl")

for json_line in loader:
    result = pipeline.process(json_line)
    saver.write(result)

集合处理

from jsonflow.core import Pipeline
from jsonflow.operators.json_ops import JsonSplitter, JsonAggregator

# 创建一个带有集合操作的管道
pipeline = Pipeline([
    JsonSplitter(split_field="items"),  # 将单个JSON拆分为多个
    TextNormalizer(),                   # 处理拆分后的每个项目
    ModelInvoker(model="gpt-3.5-turbo"),
], collection_mode="flatten")           # 自动展平结果

# 或者保持嵌套结构,并最终聚合
nested_pipeline = Pipeline([
    JsonSplitter(split_field="items"),
    TextNormalizer(),
    ModelInvoker(model="gpt-3.5-turbo"),
    JsonAggregator(aggregate_field="processed_items")
], collection_mode="nested")            # 保持嵌套结构

字段透传与系统字段

from jsonflow.core import Pipeline
from jsonflow.operators.json_ops import IdAdder, TimestampAdder

# 设置带有系统字段和透传功能的管道
pipeline = Pipeline([
    IdAdder(),                          # 添加唯一ID
    TimestampAdder(),                   # 添加时间戳
    TextNormalizer(),
    ModelInvoker(model="gpt-3.5-turbo"),
])

# 设置需要透传的字段
pipeline.set_passthrough_fields(['id', 'timestamp'])

并发处理

from jsonflow.core import Pipeline, MultiThreadExecutor
from jsonflow.io import JsonLoader, JsonSaver

# 创建多线程执行器
pipeline = Pipeline([...])
executor = MultiThreadExecutor(pipeline, max_workers=4)

# 并行处理所有数据,保持原始顺序
loader = JsonLoader("input.jsonl")
json_data_list = loader.load()
results = executor.execute_all(json_data_list)

saver = JsonSaver("output.jsonl")
saver.write_all(results)

进阶用法

表达式操作符

from jsonflow.operators.json_ops import JsonExpressionOperator

expr_op = JsonExpressionOperator({
    # 使用Lambda函数计算
    "total": lambda d: sum(item["price"] for item in d["items"]),
    # 提取和格式化字段
    "summary": lambda d: f"{d['user']['name']}的订单总额为{d['total']}元",
})

自定义操作符

from jsonflow.core import JsonOperator

class MyOperator(JsonOperator):
    def __init__(self, param, name=None):
        super().__init__(name, "My custom operator")
        self.param = param
    
    def process_item(self, json_data):
        # 处理逻辑
        json_data["result"] = json_data.get("input", "") + self.param
        return json_data

异步执行

from jsonflow.core import AsyncExecutor

# 创建异步执行器
async_executor = AsyncExecutor(pipeline)
result = await async_executor.execute(json_data)

更多示例

查看 examples 目录获取更多使用示例和最佳实践。

贡献

欢迎提交 Issue 和 Pull Request,一起完善这个库!

JSONL检查工具

这个脚本用于检查JSONL文件,验证每行是否是有效的JSON,并可以选择过滤无效行,只输出有效的JSON行。

功能特点

  • 检查JSONL文件中每行是否是有效的JSON
  • 提供选项过滤无效JSON行
  • 支持从标准输入读取和向标准输出写入
  • 提供详细的错误报告
  • 提供简单的统计信息
  • 尝试自动修复常见的JSON错误
  • 规范化空白字符处理

使用方法

./check_jsonl.py [-h] [-o OUTPUT] [-r] [-v] [-c] [-n] [-f] input

参数说明

  • input: 输入JSONL文件 (使用 "-" 从标准输入读取)
  • -o, --output OUTPUT: 输出文件 (使用 "-" 输出到标准输出)
  • -r, --remove-invalid: 移除无效的JSON行
  • -v, --verbose: 显示详细信息
  • -c, --count-only: 仅显示统计信息
  • -n, --normalize-whitespace: 规范化空白字符(将制表符、回车等替换为空格)
  • -f, --fix-errors: 尝试修复简单的JSON错误
  • -h, --help: 显示帮助信息

示例

  1. 检查JSONL文件并显示统计信息:
./check_jsonl.py data.jsonl -v
  1. 移除无效JSON行并输出到新文件:
./check_jsonl.py data.jsonl -r -o filtered_data.jsonl
  1. 从标准输入读取,过滤后输出到标准输出:
cat data.jsonl | ./check_jsonl.py - -r
  1. 只显示统计信息:
./check_jsonl.py data.jsonl -c
  1. 尝试修复JSON错误并保存结果:
./check_jsonl.py data.jsonl -f -r -o fixed_data.jsonl
  1. 规范化空白字符并过滤:
./check_jsonl.py data.jsonl -n -r > cleaned_data.jsonl

返回值

  • 0: 所有行都是有效的JSON
  • 1: 存在无效的JSON行

依赖

  • Python 3.6+
  • 标准库: argparse, json, sys, pathlib

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

guru4elephant_jsonflow-0.1.0.tar.gz (34.6 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

guru4elephant_jsonflow-0.1.0-py3-none-any.whl (47.1 kB view details)

Uploaded Python 3

File details

Details for the file guru4elephant_jsonflow-0.1.0.tar.gz.

File metadata

  • Download URL: guru4elephant_jsonflow-0.1.0.tar.gz
  • Upload date:
  • Size: 34.6 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.1.0 CPython/3.12.8

File hashes

Hashes for guru4elephant_jsonflow-0.1.0.tar.gz
Algorithm Hash digest
SHA256 ccf5dc8da46d0110e3008a72d121351f8aba00873fdddb2e778a9e1363613ddb
MD5 b242fed7a9a52e8b89651ad0f69c5584
BLAKE2b-256 a3358f7de1a8bce4ea46f3f042e7759473b89b0119e83310f9e16cf9e80d5455

See more details on using hashes here.

File details

Details for the file guru4elephant_jsonflow-0.1.0-py3-none-any.whl.

File metadata

File hashes

Hashes for guru4elephant_jsonflow-0.1.0-py3-none-any.whl
Algorithm Hash digest
SHA256 7361c972078da502325db1c70cf0ffb7549f2862d7c616881cf62b6a38820c12
MD5 77839082b755b07574b6d8330a9112b1
BLAKE2b-256 c9c39b429bfa735101e803d77223f4403b812094910ab102a8eb14de23de0c89

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page