# http-server **Repository Path**: dufafei/http-server ## Basic Information - **Project Name**: http-server - **Description**: 数据集成子项目 - HTTP服务端 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-30 - **Last Updated**: 2026-07-13 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # HTTP服务端 ## 概述 本项目为数据集成子项目,专为 HTTP 输入插件 设计。 基于 Netty 构建高性能网络接入层,结合 LMAX Disruptor 实现无锁化内存缓冲。 提供高吞吐、低延迟的 HTTP 数据接收能力。 > **完善中** ## 功能 ### 1. 支持Schema自动推导 系统支持基于配置的 JSON 请求体样例,自动解析数据结构并生成 Schema 模型。 规则如下: - **基础类型映射**: - 字段名称映射:直接将 JSON 的 Key 作为对应的字段名称。 - 字段类型映射:根据 JSON 的 Value 值类型,自动映射为对应的字段类型。 - `String` 映射为 `StringField` - `Boolean` 映射为 `BooleanField` - `Integer / Long / BigInteger` 映射为 `BigIntField` - `Float / Double / Decimal` 映射为 `DecimalField` - **嵌套对象扁平化**: - 字段名称映射:递归解析嵌套的 JSON 对象,将多层级路径使用下划线(`_`)拼接作为字段名。 - 字段类型映射:同基础类型映射。 - **数组结构处理**: - 字段名称映射:递归解析嵌套的 JSON 对象,将多层级路径使用下划线(`_`)拼接作为字段名。 - 字段类型映射:统一映射为`StringField`,取值的时候将数组转为字符串。 > 使用"\_"拼接,是因为函数执行引擎的变量只支持标识符和"_"的组合。 **示例**: 输入 JSON 报文: ```json { "userId": 10086, "isActive": true, "accountBalance": 99.99, "userProfile": { "firstName": "Alice", "age": 25, "contactInfo": { "email": "alice@example.com", "phone": "13800138000" } }, "tags": ["VIP", "NewUser"] } ``` 自动推导生成的 Schema 结构: ```json root |-- userId: bigint |-- isActive: boolean |-- accountBalance: decimal(38,18) |-- userProfile_firstName: string(255) |-- userProfile_age: bigint |-- userProfile_contactInfo_email: string(255) |-- userProfile_contactInfo_phone: string(255) |-- tags: string(255) ``` ### 2. 支持根据Schema构建数据流 根据指定的Schema,自动构建数据流对象。 ```java public static Row getRow(Schema schema, String jsonBody) { JSONObject jsonObject = JSONObject.parseObject(jsonBody); Map flatMap = new HashMap<>(schema.size()); flattenJsonToMap(jsonObject, "", flatMap); Row row = new Row(schema.size()); for (int i = 0; i < schema.size(); i++) { String fieldName = schema.getField(i).getName(); Object value = flatMap.get(fieldName); if (value == null) { row.setValue(i, null); } else if (value instanceof String) { row.setValue(i, StringValue.fromString((String) value)); } else if (value instanceof Boolean) { row.setValue(i, BooleanValue.fromBoolean((boolean) value)); } else if (value instanceof Integer) { row.setValue(i, DecimalValue.fromInt((int) value)); } else if (value instanceof Long) { row.setValue(i, DecimalValue.fromLong((long) value)); } else if (value instanceof BigInteger) { row.setValue(i, DecimalValue.fromBigInteger((BigInteger) value)); } else if (value instanceof Float) { row.setValue(i, DecimalValue.fromFloat((float) value)); } else if (value instanceof Double) { row.setValue(i, DecimalValue.fromDouble((double) value)); } else if (value instanceof BigDecimal) { row.setValue(i, DecimalValue.fromBigDecimal((BigDecimal) value)); } else { row.setValue(i, StringValue.fromString(value.toString())); } } return row; } private static void flattenJsonToMap(JSONObject json, String prefix, Map map) { for (Map.Entry entry : json.entrySet()) { String key = entry.getKey(); String fieldName = prefix.isEmpty() ? key : prefix + "_" + key; Object value = entry.getValue(); if (value instanceof JSONObject) { flattenJsonToMap((JSONObject) value, fieldName, map); } else { map.put(fieldName, value); } } } ``` ### 3. 支持同步返回和异步回调 http server 支持同步和异步返回。 http server 通过客户端请求头中是否带有 X-Async-Mode: true判断是同步还是异步返回。 #### 3.1 同步返回 #### 3.2 异步回调 异步回调需要结合回调输出插件实现,可以显著提升http输入插件的处理性能。 实现方案也分两种。 ##### 3.2.1 服务端主动推送 客户端向 ETL 系统的 HTTP Server 发送请求时,将回调地址(`CallbackURL`)一并传入。 ```json { "xx":"xx" // 任务元数据 "callbackUrl": "https://client-system.com/webhooks/etl-results" } ``` HTTP Server 收到请求后返回 202 Accepted,HTTP 连接释放。 Worker 节点完成数据处理后,发起请求将处理结果推送到客户端提供的回调地址。 ##### **3.2.2 客户端间断轮询** HTTP Server 接收请求后,生成全局唯一的 `TaskID`,HTTP Server 立即向客户端返回 `202 Accepted`,响应体包含 `TaskID` 和任务状态链接。 ```json { "taskId": "xxxxxxxx", "status": "PENDING", "message": "ETL task completed successfully.", "statusUrl": "/api/v1/tasks/${taskId}/status", // 任务状态链接, HTTP Server 生成 "resultUrl": "/api/v1/tasks/${taskId}/result" // 结果链接, HTTP Server 生成 } ``` 任务提交时:HTTP Server 在 Redis 中写入 status = PENDING。 任务执行时:Worker 节点开始处理时,更新 Redis 中的状态为 status = RUNNING。 任务完成时:Worker 节点处理完毕后,更新 Redis 中的状态为 status = SUCCESS(或 FAILED)。 FAILED下通常会附带错误信息(Error Message)。 客户端在调用任务状态链接,确认任务状态为 SUCCESS 后,再调用结果链接。 平台收到客户端的拉取请求后,通过 taskId 去存储中查询结果并返回。 > 补充说明:TaskID需要定义清理策略,自动清理 N 天(如 1 天)前的历史任务记录。 #### 3.3 提供方式 数据查询: 1.直接内存返回 数据下载: 1.平台边读边通过 HTTP Response 输出流写给客户端,虽然内存占用极低,但会长时间占用 HTTP 连接。 2.HTTP Server 生成一个带有过期时间和签名的临时下载链接,通过 HTTP `302 Redirect` 重定向给客户端,客户端下载。 ## 使用 ## 声明 本项目代码受版权保护。除明确授权外,**保留所有权利(All Rights Reserved)**。 - **授权范围**:本项目仅供内部学习、测试及非商业性质的个人研究使用。 - **禁止事项**:未经授权,严禁任何形式的二次开源、分发或作为衍生作品发布于公共代码托管平台。