Compare commits

...
12 Commits
Author SHA1 Message Date
ygxbnet 59cec93734 refactor(uploader): 优化文件上传并发处理逻辑
构建上传工具 / build (push) Successful in 1m20s
- 将循环方式从索引遍历改为range方式
- 添加处理进程创建注释说明
- 在进度更新后添加延时避免过度频繁更新
- 修复空行处理时计数器递增逻辑
- 优化panic恢复处理流程
2026-07-28 09:17:50 +08:00
ygxbnet c7755b7c29 config: 更新默认服务器地址配置
构建上传工具 / build (push) Successful in 1m14s
- 注释掉本地回环地址配置
- 修改默认URL为生产环境地址 http://112.124.71.39:8080
- 保留其他配置项不变
2026-07-27 00:55:50 +08:00
ygxbnet 2c0207c45f refactor(api): 优化抖音API请求重试机制和客户端配置
构建上传工具 / build (push) Successful in 1m12s
- 实现了抖音API请求失败时的自动重试逻辑,间隔100毫秒
- 添加了上下文取消支持,避免无限重试
- 将HTTP客户端改为全局变量,提升性能和连接复用
- 前端上传线程数限制从100调整为500
- 更新类型声明以符合最新Go语言规范
2026-07-27 00:41:14 +08:00
ygxbnet 8889e9ab55 feat(api): 根据用户是否有视频动态切换不同的token
构建上传工具 / build (push) Successful in 1m26s
- 添加了 gjson 库用于解析 JSON 数据
- 修改循环语法从传统的 for i := 0; i < n; i++ 到 range 语法
- 实现了 CheckHasVideo 函数来检查用户是否有视频作品
- 根据用户是否有视频动态选择 HasVideoToken 或 NoVideoToken
- 添加了新的 HTTP 请求头配置用于 API 调用
- 更新了依赖包并添加了相关间接依赖
2026-07-26 22:41:34 +08:00
ygxbnet 70c903f876 feat(config): 添加两个令牌配置选项
构建上传工具 / build (push) Successful in 3m26s
2026-07-26 18:13:16 +08:00
ygxbnet b0ee886b75 feat(frontend): 增加两个token的输入框
构建上传工具 / build (push) Failing after 2m46s
2026-07-26 18:09:59 +08:00
ygxbnet e15debc962 fix(config): 修复配置目录选择后未保存的问题
构建上传工具 / build (push) Successful in 2m34s
- 在用户选择目录后调用 writeCheckDir() 方法保存配置
- 确保选中的目录路径能够正确写入到配置文件中
2026-06-29 23:14:45 +08:00
ygxbnet 875a6858fe fix(uploader): 修复文件清空逻辑并优化上传流程
构建上传工具 / build (push) Successful in 2m12s
- 移除无需提示直接清空文件的功能,统一通过用户确认处理
- 修正条件判断逻辑,将 ClearFilesNoPrompt 的取反操作改为显式比较
- 调整代码结构,将文件复制到tmp目录的逻辑移到确认处理之后
- 更新取消上传的处理方式,当用户未确认时直接返回并提示等待
- 简化文件复制和清空的执行流程,提高代码可读性
2026-06-29 17:06:46 +08:00
ygxbnet 94fe4f8f62 添加注释
构建上传工具 / build (push) Successful in 4m29s
2026-06-29 16:55:15 +08:00
ygxbnet b100fc9b32 fix(upload): 修正上传文件清空确认对话框文本
构建上传工具 / build (push) Successful in 1m51s
- 修正了标题文本"是否确认清空需要上传文件"
- 修正了提示文本中的"tmp"文件夹显示格式
- 更新了确认对话框的内容描述文本
2026-05-01 22:51:12 +08:00
ygxbnet d426c16104 refactor(uploader): 优化上传器代码结构和上下文取消处理
- 在进度更新循环中添加上下文取消检查点
- 在文件复制操作前添加上下文取消检查点
- 重构代码缩进和括号位置以提高可读性
- 优化 goroutine 中的上下文取消处理逻辑
- 统一代码块的括号格式和缩进风格
2026-05-01 22:47:39 +08:00
ygxbnet 993814cdfa refactor(uploader): 优化文件信息统计逻辑
- 将变量名 fInfo 重命名为 filesInfo 以提高可读性
- 调整代码顺序,将 AddLog 调用移到变量声明后
- 统一使用新变量名在所有相关位置进行引用
- 移动 g.SetLimit 注释位置以提高代码可读性
2026-05-01 22:37:45 +08:00
10 changed files with 249 additions and 135 deletions
+22 -9
View File
@@ -9,7 +9,8 @@ import {configModel} from "@/model.ts";
import Config = config.Config; import Config = config.Config;
// const serverUrl = ref('') // const serverUrl = ref('')
const token = ref('') const hasVideoToken = ref('')
const noVideoToken = ref('')
const checkDir = ref('') const checkDir = ref('')
const concurrentFiles = ref(1) const concurrentFiles = ref(1)
const uploadThreads = ref(1) const uploadThreads = ref(1)
@@ -59,6 +60,7 @@ const selectDirectory = () => {
SelectPath().then((path) => { SelectPath().then((path) => {
if (path) { if (path) {
checkDir.value = path checkDir.value = path
writeCheckDir()
} else { } else {
ElMessage.warning('未选择目录,不更改配置') ElMessage.warning('未选择目录,不更改配置')
} }
@@ -70,8 +72,8 @@ const startRun = () => {
// ElMessage.warning('请输入服务器地址') // ElMessage.warning('请输入服务器地址')
// return // return
// } // }
if (!token.value) { if (!hasVideoToken.value || !noVideoToken.value) {
ElMessage.error('请输入Token') ElMessage.error('请输入两个上传 Token')
return return
} }
if (!checkDir.value) { if (!checkDir.value) {
@@ -104,8 +106,11 @@ const clearLog = () => {
// const writeServerUrl =() => { // const writeServerUrl =() => {
// WriteConfig("url", serverUrl.value) // WriteConfig("url", serverUrl.value)
// } // }
const writeToken = () => { const writeHasVideoToken = () => {
WriteConfig(configModel.Token, token.value) WriteConfig(configModel.HasVideoToken, hasVideoToken.value)
}
const writeNoVideoToken = () => {
WriteConfig(configModel.NoVideoToken, noVideoToken.value)
} }
const writeCheckDir = () => { const writeCheckDir = () => {
WriteConfig(configModel.CheckDir, checkDir.value) WriteConfig(configModel.CheckDir, checkDir.value)
@@ -124,7 +129,8 @@ const writeUploadThreads = () => {
try { try {
GetConfig().then((config: Config) => { GetConfig().then((config: Config) => {
// serverUrl.value = config.url // serverUrl.value = config.url
token.value = config.token hasVideoToken.value = config.has_video_token
noVideoToken.value = config.no_video_token
checkDir.value = config.check_dir checkDir.value = config.check_dir
concurrentFiles.value = config.handle_file_count concurrentFiles.value = config.handle_file_count
uploadThreads.value = config.thread_count uploadThreads.value = config.thread_count
@@ -167,8 +173,15 @@ try {
<!-- </div>--> <!-- </div>-->
<div class="form-item"> <div class="form-item">
<label>Token</label> <label>有视频 Token</label>
<el-input v-model="token" placeholder="请输入Token" :disabled="isRunning" @change="writeToken()"/> <el-input v-model="hasVideoToken" placeholder="有视频账号上传使用的 Token" :disabled="isRunning"
@change="writeHasVideoToken()"/>
</div>
<div class="form-item">
<label>无视频 Token</label>
<el-input v-model="noVideoToken" placeholder="无视频账号上传使用的 Token" :disabled="isRunning"
@change="writeNoVideoToken()"/>
</div> </div>
<div class="form-item"> <div class="form-item">
@@ -187,7 +200,7 @@ try {
<div class="form-item"> <div class="form-item">
<label>单文件上传线程</label> <label>单文件上传线程</label>
<el-input-number v-model="uploadThreads" :min="1" :max="100" :disabled="isRunning" <el-input-number v-model="uploadThreads" :min="1" :max="500" :disabled="isRunning"
@change="writeUploadThreads()"/> @change="writeUploadThreads()"/>
</div> </div>
@@ -40,9 +40,9 @@ const handleConfirm = () => {
align-center align-center
> >
<template #header> <template #header>
是否确认清空上传文件 是否确认清空需要上传文件
</template> </template>
<div class="hint-text">以下文件将会被清空并移动到tmp文件夹进行上传,您是否确认</div> <div class="hint-text">以下文件将会被清空并移动到 tmp 文件夹进行上传,您是否确认</div>
<div class="dialog-content"> <div class="dialog-content">
<div class="file-list"> <div class="file-list">
<div v-for="(file, index) in fileList" :key="index" class="file-item"> <div v-for="(file, index) in fileList" :key="index" class="file-item">
+2
View File
@@ -1,6 +1,8 @@
export enum configModel { export enum configModel {
Url = "url", Url = "url",
Token = "token", Token = "token",
HasVideoToken = "has-video-token",
NoVideoToken = "no-video-token",
ThreadCount = "thread-count", ThreadCount = "thread-count",
HandleFileCount = "handle-file-count", HandleFileCount = "handle-file-count",
IsRunOnStart = "is-run-on-start", IsRunOnStart = "is-run-on-start",
+4
View File
@@ -3,6 +3,8 @@ export namespace config {
export class Config { export class Config {
url: string; url: string;
token: string; token: string;
has_video_token: string;
no_video_token: string;
thread_count: number; thread_count: number;
handle_file_count: number; handle_file_count: number;
is_run_on_start: boolean; is_run_on_start: boolean;
@@ -17,6 +19,8 @@ export namespace config {
if ('string' === typeof source) source = JSON.parse(source); if ('string' === typeof source) source = JSON.parse(source);
this.url = source["url"]; this.url = source["url"];
this.token = source["token"]; this.token = source["token"];
this.has_video_token = source["has_video_token"];
this.no_video_token = source["no_video_token"];
this.thread_count = source["thread_count"]; this.thread_count = source["thread_count"];
this.handle_file_count = source["handle_file_count"]; this.handle_file_count = source["handle_file_count"];
this.is_run_on_start = source["is_run_on_start"]; this.is_run_on_start = source["is_run_on_start"];
+3
View File
@@ -5,6 +5,7 @@ go 1.26
require ( require (
github.com/fsnotify/fsnotify v1.9.0 github.com/fsnotify/fsnotify v1.9.0
github.com/spf13/viper v1.21.0 github.com/spf13/viper v1.21.0
github.com/tidwall/gjson v1.14.2
github.com/wailsapp/wails/v2 v2.12.0 github.com/wailsapp/wails/v2 v2.12.0
golang.org/x/sync v0.20.0 golang.org/x/sync v0.20.0
) )
@@ -37,6 +38,8 @@ require (
github.com/spf13/cast v1.10.0 // indirect github.com/spf13/cast v1.10.0 // indirect
github.com/spf13/pflag v1.0.10 // indirect github.com/spf13/pflag v1.0.10 // indirect
github.com/subosito/gotenv v1.6.0 // indirect github.com/subosito/gotenv v1.6.0 // indirect
github.com/tidwall/match v1.1.1 // indirect
github.com/tidwall/pretty v1.2.0 // indirect
github.com/tkrajina/go-reflector v0.5.8 // indirect github.com/tkrajina/go-reflector v0.5.8 // indirect
github.com/valyala/bytebufferpool v1.0.0 // indirect github.com/valyala/bytebufferpool v1.0.0 // indirect
github.com/valyala/fasttemplate v1.2.2 // indirect github.com/valyala/fasttemplate v1.2.2 // indirect
+6
View File
@@ -81,6 +81,12 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8=
github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU=
github.com/tidwall/gjson v1.14.2 h1:6BBkirS0rAHjumnjHF6qgy5d2YAJ1TLIaFE2lzfOLqo=
github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk=
github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA=
github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM=
github.com/tidwall/pretty v1.2.0 h1:RWIZEg2iJ8/g6fDDYzMpobmaoGh5OLl4AXtGUGPcqCs=
github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU=
github.com/tkrajina/go-reflector v0.5.8 h1:yPADHrwmUbMq4RGEyaOUpz2H90sRsETNVpjzo3DLVQQ= github.com/tkrajina/go-reflector v0.5.8 h1:yPADHrwmUbMq4RGEyaOUpz2H90sRsETNVpjzo3DLVQQ=
github.com/tkrajina/go-reflector v0.5.8/go.mod h1:ECbqLgccecY5kPmPmXg1MrHW585yMcDkVl6IvJe64T4= github.com/tkrajina/go-reflector v0.5.8/go.mod h1:ECbqLgccecY5kPmPmXg1MrHW585yMcDkVl6IvJe64T4=
github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw=
+67 -3
View File
@@ -7,8 +7,11 @@ import (
"io" "io"
"net/http" "net/http"
"net/url" "net/url"
"strings"
"sync" "sync"
"time" "time"
"github.com/tidwall/gjson"
) )
var httpClient = &http.Client{ var httpClient = &http.Client{
@@ -31,9 +34,9 @@ func init() {
func InitConn() { func InitConn() {
wg := &sync.WaitGroup{} wg := &sync.WaitGroup{}
for i := 0; i < 10; i++ { for range 10 {
wg.Go(func() { wg.Go(func() {
for i := 0; i < 50; i++ { for range 50 {
resp, err := httpClient.Get(config.APPConfig.Url + "/api/test") resp, err := httpClient.Get(config.APPConfig.Url + "/api/test")
if err != nil { if err != nil {
fmt.Println(err) fmt.Println(err)
@@ -56,8 +59,31 @@ func UploadDataToServer(ctx context.Context, data string) error {
<-limit <-limit
}() }()
//根据账号是否有视频判断要存储的token
params := url.Values{} params := url.Values{}
params.Set("token", config.APPConfig.Token) secUserId := strings.Split(data, "----")[1]
hasVideo, err := CheckHasVideo(secUserId)
if err != nil {
for err != nil {
fmt.Println("抖音请求api重试")
//定时器定时100毫秒
ticker := time.NewTicker(100 * time.Millisecond)
select {
case <-ctx.Done():
return err
case <-ticker.C:
ticker.Stop()
hasVideo, err = CheckHasVideo(secUserId)
}
}
}
if hasVideo {
params.Set("token", config.APPConfig.HasVideoToken)
} else {
params.Set("token", config.APPConfig.NoVideoToken)
}
params.Set("data", data) params.Set("data", data)
//http://127.0.0.1:8080/api/data?token=123456&data=123456 //http://127.0.0.1:8080/api/data?token=123456&data=123456
@@ -82,3 +108,41 @@ func UploadDataToServer(ctx context.Context, data string) error {
return err return err
} }
var douYinClient = &http.Client{}
func CheckHasVideo(secUserId string) (bool, error) {
req, err := http.NewRequest(
"GET",
"https://imdesktop.douyin.com/aweme/v1/web/user/profile/other/?sec_user_id="+secUserId,
nil,
)
if err != nil {
fmt.Println(err)
return false, err
}
req.Header.Add("User-Agent", "Apifox/1.0.0 (https://apifox.com)")
req.Header.Add("Accept", "*/*")
req.Header.Add("Host", "imdesktop.douyin.com")
req.Header.Add("Connection", "keep-alive")
res, err := douYinClient.Do(req)
if err != nil {
fmt.Println(err)
return false, err
}
defer res.Body.Close()
body, err := io.ReadAll(res.Body)
if err != nil {
fmt.Println(err)
return false, err
}
if gjson.Get(string(body), "user.aweme_count").Int() > 0 {
return true, nil
}
return false, nil
}
+10 -1
View File
@@ -11,6 +11,8 @@ import (
type Config struct { type Config struct {
Url string `json:"url" mapstructure:"url"` Url string `json:"url" mapstructure:"url"`
Token string `json:"token" mapstructure:"token"` Token string `json:"token" mapstructure:"token"`
HasVideoToken string `json:"has_video_token" mapstructure:"has-video-token"`
NoVideoToken string `json:"no_video_token" mapstructure:"no-video-token"`
ThreadCount int `json:"thread_count" mapstructure:"thread-count"` ThreadCount int `json:"thread_count" mapstructure:"thread-count"`
HandleFileCount int `json:"handle_file_count" mapstructure:"handle-file-count"` HandleFileCount int `json:"handle_file_count" mapstructure:"handle-file-count"`
IsRunOnStart bool `json:"is_run_on_start" mapstructure:"is-run-on-start"` IsRunOnStart bool `json:"is_run_on_start" mapstructure:"is-run-on-start"`
@@ -24,6 +26,8 @@ var configMu sync.Mutex
const ( const (
Url = "url" Url = "url"
Token = "token" Token = "token"
HasVideoToken = "has-video-token"
NoVideoToken = "no-video-token"
ThreadCount = "thread-count" ThreadCount = "thread-count"
HandleFileCount = "handle-file-count" HandleFileCount = "handle-file-count"
IsRunOnStart = "is-run-on-start" IsRunOnStart = "is-run-on-start"
@@ -34,8 +38,11 @@ const (
func InitConfig() { func InitConfig() {
// 设置默认配置 // 设置默认配置
defaultConfig := Config{ defaultConfig := Config{
Url: "http://127.0.0.1:8080", //Url: "http://127.0.0.1:8080",
Url: "http://112.124.71.39:8080",
Token: "", Token: "",
HasVideoToken: "1234",
NoVideoToken: "5678",
ThreadCount: 10, ThreadCount: 10,
HandleFileCount: 25, HandleFileCount: 25,
IsRunOnStart: false, IsRunOnStart: false,
@@ -44,6 +51,8 @@ func InitConfig() {
} }
viper.SetDefault(Url, defaultConfig.Url) viper.SetDefault(Url, defaultConfig.Url)
viper.SetDefault(Token, defaultConfig.Token) viper.SetDefault(Token, defaultConfig.Token)
viper.SetDefault(HasVideoToken, defaultConfig.HasVideoToken)
viper.SetDefault(NoVideoToken, defaultConfig.NoVideoToken)
viper.SetDefault(ThreadCount, defaultConfig.ThreadCount) viper.SetDefault(ThreadCount, defaultConfig.ThreadCount)
viper.SetDefault(HandleFileCount, defaultConfig.HandleFileCount) viper.SetDefault(HandleFileCount, defaultConfig.HandleFileCount)
viper.SetDefault(IsRunOnStart, defaultConfig.IsRunOnStart) viper.SetDefault(IsRunOnStart, defaultConfig.IsRunOnStart)
+53 -40
View File
@@ -44,10 +44,12 @@ func StartUpload(ctx context.Context, logChan *chan string) {
AddLog(logChan, `单文件上传线程: `+strconv.Itoa(config.APPConfig.ThreadCount)) AddLog(logChan, `单文件上传线程: `+strconv.Itoa(config.APPConfig.ThreadCount))
AddLog(logChan, "===============================================") AddLog(logChan, "===============================================")
//创建连接池
AddLog(logChan, "正在创建连接池(连接池可避免首次大量上传时出现网络错误)") AddLog(logChan, "正在创建连接池(连接池可避免首次大量上传时出现网络错误)")
api.InitConn() api.InitConn()
AddLog(logChan, "创建连接池完成,开始运行程序") AddLog(logChan, "创建连接池完成,开始运行程序")
//清除进度
progress.Clear() progress.Clear()
//推送上传进度 //推送上传进度
@@ -57,6 +59,8 @@ func StartUpload(ctx context.Context, logChan *chan string) {
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
}
var pg []Progress var pg []Progress
progress.Range(func(_, value any) bool { progress.Range(func(_, value any) bool {
pg = append(pg, value.(Progress)) pg = append(pg, value.(Progress))
@@ -66,12 +70,13 @@ func StartUpload(ctx context.Context, logChan *chan string) {
time.Sleep(250 * time.Millisecond) time.Sleep(250 * time.Millisecond)
} }
}
}() }()
//开启上传程序 //开启上传程序
for { for {
uploadData(ctx, logChan) uploadData(ctx, logChan)
//延时1分钟运行
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
@@ -116,22 +121,8 @@ func uploadData(ctx context.Context, logChan *chan string) {
return return
} }
//是否向用户提示清空文件,并复制文件到tmp //向用户提示清空文件,获取确认
if config.APPConfig.ClearFilesNoPrompt { if config.APPConfig.ClearFilesNoPrompt == false {
//不用提示直接复制文件到tmp
for _, p := range f {
err := copyFile(p, "./tmp/"+filepath.Base(p))
if err != nil {
AddLog(logChan, "复制文件失败:"+err.Error())
} else {
files = append(files, "./tmp/"+filepath.Base(p))
err := os.Truncate(p, 0)
if err != nil {
AddLog(logChan, "清空文件失败:"+err.Error())
}
}
}
} else {
//提示用户 //提示用户
wailsruntime.EventsEmit(ctx, "clear-files", f) wailsruntime.EventsEmit(ctx, "clear-files", f)
@@ -141,8 +132,21 @@ func uploadData(ctx context.Context, logChan *chan string) {
confirm <- optionalData[0].(bool) confirm <- optionalData[0].(bool)
}) })
if <-confirm { if <-confirm == false {
//取消上传
AddLog(logChan, "已取消上传,1分钟后再运行")
return
}
}
//复制文件到tmp目录
for _, p := range f { for _, p := range f {
select {
case <-ctx.Done():
return
default:
}
err := copyFile(p, "./tmp/"+filepath.Base(p)) err := copyFile(p, "./tmp/"+filepath.Base(p))
if err != nil { if err != nil {
AddLog(logChan, "复制文件失败:"+err.Error()) AddLog(logChan, "复制文件失败:"+err.Error())
@@ -154,25 +158,22 @@ func uploadData(ctx context.Context, logChan *chan string) {
} }
} }
} }
} else {
AddLog(logChan, "已取消上传,1分钟后再运行")
return
}
}
} }
//检测到文件 //检测到文件
//统计文件行数 //统计文件行数
var fInfo = make(map[string]fileInfo) var filesInfo = make(map[string]fileInfo)
AddLog(logChan, fmt.Sprintf("正在统计 %v 个文件行数", len(files)))
isAllEmpty := true isAllEmpty := true
AddLog(logChan, fmt.Sprintf("正在统计 %v 个文件行数", len(files)))
for _, filePath := range files { for _, filePath := range files {
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
}
file, err := os.Open(filePath) file, err := os.Open(filePath)
if err != nil { if err != nil {
AddLog(logChan, "打开文件失败:"+err.Error()) AddLog(logChan, "打开文件失败:"+err.Error())
@@ -189,7 +190,7 @@ func uploadData(ctx context.Context, logChan *chan string) {
if lineCount == 0 { if lineCount == 0 {
continue continue
} }
fInfo[filepath.Base(filePath)] = fileInfo{ filesInfo[filepath.Base(filePath)] = fileInfo{
FilePath: filePath, FilePath: filePath,
FileLines: lineCount, FileLines: lineCount,
} }
@@ -198,7 +199,6 @@ func uploadData(ctx context.Context, logChan *chan string) {
AddLog(logChan, fmt.Sprintf("%s 文件行数:%v", filepath.Base(filePath), lineCount)) AddLog(logChan, fmt.Sprintf("%s 文件行数:%v", filepath.Base(filePath), lineCount))
} }
}
if isAllEmpty { if isAllEmpty {
AddLog(logChan, "所有文件都为空,不进行上传") AddLog(logChan, "所有文件都为空,不进行上传")
@@ -207,7 +207,7 @@ func uploadData(ctx context.Context, logChan *chan string) {
//刷新文件上传进度 //刷新文件上传进度
progress.Clear() progress.Clear()
for fileName, info := range fInfo { for fileName, info := range filesInfo {
progress.Store(fileName, progress.Store(fileName,
Progress{ Progress{
FileName: fileName, FileName: fileName,
@@ -220,18 +220,23 @@ func uploadData(ctx context.Context, logChan *chan string) {
// 使用 errgroup 控制同时处理的文件数,并开始上传文件任务 // 使用 errgroup 控制同时处理的文件数,并开始上传文件任务
g, egctx := errgroup.WithContext(ctx) g, egctx := errgroup.WithContext(ctx)
g.SetLimit(config.APPConfig.HandleFileCount) // 设置同时处理文件数 // 设置同时处理文件数
g.SetLimit(config.APPConfig.HandleFileCount)
// 执行文件上传任务参数(文件路径,文件行数) // 执行文件上传任务参数(文件路径,文件行数)
for fileName, info := range fInfo { for fileName, info := range filesInfo {
select { select {
case <-egctx.Done(): case <-egctx.Done():
return return
default: default:
}
g.Go(func() error { g.Go(func() error {
select { select {
case <-egctx.Done(): case <-egctx.Done():
return egctx.Err() return egctx.Err()
default: default:
}
AddLog(logChan, "正在上传文件:"+fileName) AddLog(logChan, "正在上传文件:"+fileName)
processFile(egctx, logChan, info.FilePath, info.FileLines) processFile(egctx, logChan, info.FilePath, info.FileLines)
@@ -240,28 +245,28 @@ func uploadData(ctx context.Context, logChan *chan string) {
case <-egctx.Done(): case <-egctx.Done():
return egctx.Err() return egctx.Err()
default: default:
}
//上传完成,删除缓存文件 //上传完成,删除缓存文件
err := os.Remove(info.FilePath) err := os.Remove(info.FilePath)
if err != nil { if err != nil {
AddLog(logChan, "删除缓存文件失败:"+err.Error()) AddLog(logChan, "删除缓存文件失败:"+err.Error())
} }
return nil return nil
}
}
}) })
} }
}
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
}
// 等待所有任务完成 // 等待所有任务完成
g.Wait() g.Wait()
AddLog(logChan, "所有任务执行完成!") AddLog(logChan, "所有任务执行完成!")
AddLog(logChan, fmt.Sprintf("上传完成,耗时:%s", time.Since(start).String())) AddLog(logChan, fmt.Sprintf("上传完成,耗时:%s", time.Since(start).String()))
}
} }
// copyFile 快速拷贝文件 src -> dst // copyFile 快速拷贝文件 src -> dst
@@ -329,17 +334,19 @@ func processFile(ctx context.Context, logChan *chan string, filePath string, fil
lines := make(chan string, 200) lines := make(chan string, 200)
var countLine int32 = 0 var countLine int32 = 0
// 创建指定个worker同时处理文件上传 // 创建指定个worker同时处理文件上传
for i := 0; i < config.APPConfig.ThreadCount; i++ { for i := range config.APPConfig.ThreadCount {
select { select {
case <-ctx.Done(): case <-ctx.Done():
close(lines) close(lines)
return return
default: default:
}
//创建处理进程
go func() { go func() {
processLines(ctx, logChan, &lines, i, filePath, &countLine) processLines(ctx, logChan, &lines, i, filePath, &countLine)
}() }()
} }
}
// 读取文件并发送到通道 // 读取文件并发送到通道
scanner := bufio.NewScanner(file) scanner := bufio.NewScanner(file)
@@ -350,13 +357,15 @@ func processFile(ctx context.Context, logChan *chan string, filePath string, fil
fmt.Println("panic:", f+":"+strconv.Itoa(l), r) fmt.Println("panic:", f+":"+strconv.Itoa(l), r)
} }
}() }()
for scanner.Scan() { for scanner.Scan() {
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
lines <- scanner.Text()
} }
lines <- scanner.Text()
} }
}() }()
@@ -367,6 +376,8 @@ func processFile(ctx context.Context, logChan *chan string, filePath string, fil
close(lines) //关闭processLines中的上传线程 close(lines) //关闭processLines中的上传线程
return return
default: default:
}
progress.Store(filepath.Base(filePath), progress.Store(filepath.Base(filePath),
Progress{ Progress{
FileName: filepath.Base(filePath), FileName: filepath.Base(filePath),
@@ -375,7 +386,7 @@ func processFile(ctx context.Context, logChan *chan string, filePath string, fil
Percentage: int(float64(countLine) / float64(fileLines) * 100), Percentage: int(float64(countLine) / float64(fileLines) * 100),
}, },
) )
} time.Sleep(50 * time.Millisecond)
} }
//上传完成,进度设为100 //上传完成,进度设为100
progress.Store(filepath.Base(filePath), progress.Store(filepath.Base(filePath),
@@ -403,8 +414,11 @@ func processLines(ctx context.Context, logChan *chan string, lines *chan string,
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
}
// 跳过空行 // 跳过空行
if strings.TrimSpace(line) == "" { if strings.TrimSpace(line) == "" {
atomic.AddInt32(countLine, 1)
continue continue
} }
// 上传数据 // 上传数据
@@ -413,7 +427,6 @@ func processLines(ctx context.Context, logChan *chan string, lines *chan string,
} }
atomic.AddInt32(countLine, 1) atomic.AddInt32(countLine, 1)
} }
}
} }
// AddLog 添加日志 // AddLog 添加日志
+1 -1
View File
@@ -30,7 +30,7 @@ func main() {
}, },
BackgroundColour: &options.RGBA{R: 27, G: 38, B: 54, A: 1}, BackgroundColour: &options.RGBA{R: 27, G: 38, B: 54, A: 1},
OnStartup: app.startup, OnStartup: app.startup,
Bind: []interface{}{ Bind: []any{
app, app,
}, },
}) })