如何基于Java Spring Boot创建透传JSON行流的REST API
Spring Boot实现流式透传数据库JSON行并做行处理
核心方案
通过数据库流式获取结果集,结合Spring Boot的流式响应能力,在不加载全量数据到内存的前提下,对每一行JSON做过滤/转换,同时支持少量行缓存,最终将处理后的数据流实时返回给客户端。
一、数据库端流式获取JSON行
不同数据库的流式查询配置略有差异,以下以JDBC为例实现流式结果集获取:
@Autowired private DataSource dataSource; public Stream<String> getDatabaseJsonStream(String modifiedQuery) throws SQLException { try (Connection conn = dataSource.getConnection(); Statement stmt = conn.createStatement()) { // 针对MySQL开启流式获取(PostgreSQL可设置合理fetchSize如50) stmt.setFetchSize(Integer.MIN_VALUE); ResultSet rs = stmt.executeQuery(modifiedQuery); // 将ResultSet转为JSON行的Stream,自动处理资源释放 return Stream.generate(() -> { try { if (rs.next()) { // 假设数据库返回的是JSON格式字段,或自行将行转为JSON字符串 return rs.getString("json_content"); } else { rs.close(); return null; } } catch (SQLException e) { throw new RuntimeException("读取数据库流失败", e); } }).takeWhile(Objects::nonNull); } }
注意:MySQL JDBC URL需添加
useCursorFetch=true参数,确保流式生效;PostgreSQL无需特殊URL参数,设置stmt.setFetchSize(50)即可分批拉取数据。
二、Spring Boot API流式返回处理后的数据流
1. Servlet栈实现(同步异步兼容)
使用StreamingResponseBody实现流式响应,同时嵌入行过滤/转换逻辑,支持少量缓存:
@GetMapping("/stream-processed-data") public ResponseEntity<StreamingResponseBody> streamProcessedData(@RequestParam String inputParam) { // 根据请求参数生成修改后的数据库查询语句 String dbQuery = buildModifiedQuery(inputParam); try { Stream<String> rawJsonStream = getDatabaseJsonStream(dbQuery); // 处理流:过滤+转换+可选缓存(此处缓存5行) Stream<String> processedStream = rawJsonStream .filter(this::filterJsonRow) .map(this::transformJsonRow) .buffer(5) // 缓存5行批量处理,可根据业务调整 .flatMap(List::stream); // 构建流式响应体 StreamingResponseBody responseBody = outputStream -> { try (PrintWriter writer = new PrintWriter(outputStream)) { processedStream.forEach(writer::println); writer.flush(); } }; return ResponseEntity.ok() .contentType(MediaType.APPLICATION_NDJSON) // 用换行分隔JSON,适配流式解析 .body(responseBody); } catch (SQLException e) { return ResponseEntity.internalServerError().build(); } } // 根据请求参数生成数据库查询的逻辑 private String buildModifiedQuery(String inputParam) { return String.format("SELECT json_content FROM target_table WHERE filter_col = '%s'", inputParam); } // JSON行过滤逻辑示例 private boolean filterJsonRow(String jsonRow) { JSONObject json = new JSONObject(jsonRow); return json.getInt("data_status") == 1; // 只保留状态为1的行 } // JSON行转换逻辑示例 private String transformJsonRow(String jsonRow) { JSONObject json = new JSONObject(jsonRow); json.put("processed_timestamp", System.currentTimeMillis()); json.remove("sensitive_field"); // 移除敏感字段 return json.toString(); }
2. WebFlux响应式实现
若使用Spring WebFlux,可直接用Flux处理流式数据,结合R2DBC实现数据库流式查询:
@Autowired private DatabaseClient r2dbcClient; @GetMapping("/stream-reactive-data") public Flux<String> streamReactiveData(@RequestParam String inputParam) { String dbQuery = buildModifiedQuery(inputParam); return r2dbcClient.sql(dbQuery) .map(row -> row.get("json_content", String.class)) .filter(this::filterJsonRow) .map(this::transformJsonRow) .buffer(5) // 可选缓存 .flatMapIterable(Function.identity()); }
关键注意事项
- 资源释放:必须通过try-with-resources或Spring自动管理确保数据库连接、ResultSet等资源及时释放,避免连接泄漏。
- 内容类型:推荐使用
application/x-ndjson(换行分隔JSON),客户端可逐行解析,避免等待全量数据。 - 缓存大小:
buffer(n)的n值需根据业务平衡内存占用和处理效率,不宜过大导致内存压力。 - 异常处理:需为数据库流读取、行处理添加异常捕获,避免流式响应中途中断。
内容的提问来源于stack exchange,提问作者mgerbracht
相关产品推荐
相关产品推荐

