You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Java Socket多线程文件传输:如何有序接收分块文件并写入

实现多线程文件分块传输与有序合并

服务器端实现逻辑

核心步骤

  • 计算分块参数:根据文件总大小,拆分出4个分块的起始位置和大小(最后一个分块可能小于其他分块)
  • 多线程发送:每个线程负责一个分块的传输,发送前先传递分块的索引和大小,让客户端明确分块顺序

代码示例(Python)

import os
import socket
import threading

SERVER_HOST = '0.0.0.0'
SERVER_PORT_BASE = 8080
TARGET_FILE = 'your_large_file.ext'
CHUNK_NUM = 4

def send_single_chunk(chunk_idx, start_pos, chunk_size):
    # 为每个分块创建独立socket连接
    sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    sock.bind((SERVER_HOST, SERVER_PORT_BASE + chunk_idx))
    sock.listen(1)
    conn, _ = sock.accept()

    # 先发送分块元数据:索引+大小
    conn.sendall(f"{chunk_idx}:{chunk_size}".encode())
    # 读取对应分块内容并发送
    with open(TARGET_FILE, 'rb') as f:
        f.seek(start_pos)
        chunk_data = f.read(chunk_size)
        conn.sendall(chunk_data)
    
    conn.close()
    sock.close()
    print(f"分块 {chunk_idx} 发送完成")

if __name__ == '__main__':
    total_size = os.path.getsize(TARGET_FILE)
    base_chunk_size = total_size // CHUNK_NUM
    threads = []

    for idx in range(CHUNK_NUM):
        start = idx * base_chunk_size
        # 最后一个分块取剩余所有内容
        chunk_size = base_chunk_size if idx != CHUNK_NUM -1 else total_size - start
        t = threading.Thread(target=send_single_chunk, args=(idx, start, chunk_size))
        threads.append(t)
        t.start()
    
    for t in threads:
        t.join()
    print("所有分块发送完毕")

客户端实现逻辑

核心步骤

  • 多线程接收:每个线程连接对应分块的服务器端口,接收元数据和分块内容
  • 线程安全存储:用字典存储已接收分块,配合锁避免多线程写入冲突
  • 有序合并:所有分块接收完成后,按索引顺序写入磁盘

代码示例(Python)

import socket
import threading
import os

SERVER_HOST = '192.168.x.x' # 替换为服务器IP
SERVER_PORT_BASE = 8080
CHUNK_NUM = 4
OUTPUT_FILE = 'merged_file.ext'

# 存储已接收分块,键为分块索引
received_chunks = {}
chunk_lock = threading.Lock()

def receive_single_chunk(chunk_idx):
    sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    sock.connect((SERVER_HOST, SERVER_PORT_BASE + chunk_idx))

    # 接收分块元数据
    meta = sock.recv(1024).decode()
    idx, size = meta.split(':')
    idx = int(idx)
    size = int(size)

    # 循环接收分块内容,确保完整接收
    chunk_data = b''
    while len(chunk_data) < size:
        recv_buf = sock.recv(min(1024, size - len(chunk_data)))
        if not recv_buf:
            break
        chunk_data += recv_buf
    
    # 线程安全写入字典
    with chunk_lock:
        received_chunks[idx] = chunk_data
    
    sock.close()
    print(f"分块 {chunk_idx} 接收完成")

if __name__ == '__main__':
    threads = []
    for idx in range(CHUNK_NUM):
        t = threading.Thread(target=receive_single_chunk, args=(idx,))
        threads.append(t)
        t.start()
    
    for t in threads:
        t.join()
    
    # 按索引顺序合并写入文件
    with open(OUTPUT_FILE, 'wb') as f:
        for idx in range(CHUNK_NUM):
            f.write(received_chunks[idx])
    
    print("文件合并完成,已写入磁盘")

关键注意点

  • 服务器端可以优化为单端口监听,客户端连接后分配分块索引(示例用多端口是为了简化逻辑)
  • 传输大文件时,可加入校验机制(如MD5)验证分块完整性
  • 客户端的锁必须严格控制,避免多线程同时写入字典导致数据混乱

内容的提问来源于stack exchange,提问作者Vimal Arjun

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.05 08:35:28