作者 yangfu

refactor: 车间统计增加调试日志,优化统计

... ... @@ -291,12 +291,6 @@ func (productRecordService *ProductRecordService) ProductRecordStatics(cmd *comm
}()
var _ domain.ProductRecordRepository
var productRecord *domain.ProductRecord = cmd.ProductRecord
//_,productRecord,err = factory.FastPgProductRecord(transactionContext,cmd.ProductRecordId)
//if err!=nil{
// log.Logger.Error(err.Error())
// return nil, nil
//}
//
if productRecord == nil {
return nil, nil
}
... ...
package domainService
import (
"fmt"
"github.com/hibiken/asynq"
"github.com/linmadan/egglib-go/utils/json"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/constant"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/domain"
)
func SendWorkshopWorkTimeStaticJob(productRecord *domain.ProductAttendanceRecord) error {
return SendAsyncJob(domain.TaskKeyWorkshopWorkTimeRecordStatics(), productRecord)
const (
QueueProduct = "product"
QueueDevice = "device"
QueueDefault = "default"
)
func FormatQueue(qt string) string {
return fmt.Sprintf("%v:queue:%v", constant.CACHE_PREFIX, qt)
}
func SendProductRecordStaticsJob(productRecord *domain.ProductRecord) error {
task := asynq.NewTask(domain.TaskKeyPatternProductRecordStatics(), []byte(json.MarshalToString(productRecord)))
func SendWorkshopWorkTimeStaticJob(r *domain.ProductAttendanceRecord) error {
return SendAsyncJob(domain.TaskKeyWorkshopWorkTimeRecordStatics(), r, asynq.Queue(FormatQueue(QueueDefault)))
}
client := asynq.NewClient(asynq.RedisClientOpt{Addr: constant.REDIS_ADDRESS})
_, err := client.Enqueue(task)
return err
func SendProductRecordStaticsJob(r *domain.ProductRecord) error {
return SendAsyncJob(domain.TaskKeyPatternProductRecordStatics(), r, asynq.Queue(FormatQueue(QueueProduct)))
}
func SendDeviceZkTecoReportJob(productRecord *domain.DeviceZkTeco) error {
return SendAsyncJob(domain.TaskDeviceZkTecoReport(), productRecord)
func SendDeviceZkTecoReportJob(r *domain.DeviceZkTeco) error {
return SendAsyncJob(domain.TaskDeviceZkTecoReport(), r, asynq.Queue(FormatQueue(QueueDefault)))
}
func SendWorkshopDeviceData(productRecord *domain.DeviceCollection) error {
return SendAsyncJob(domain.TaskDeviceCollection(), productRecord, asynq.Queue("device"))
func SendWorkshopDeviceData(r *domain.DeviceCollection) error {
return SendAsyncJob(domain.TaskDeviceCollection(), r, asynq.Queue(FormatQueue(QueueDevice)))
}
func SendAsyncJob(queueName string, job interface{}, opts ...asynq.Option) error {
... ...
... ... @@ -27,8 +27,6 @@ func (ptr *PGProductRecordService) EmployeeProductStatics(productRecord *domain.
var (
workshopRepository, _ = repository.NewWorkshopRepository(ptr.transactionContext)
productPlanRepository, _ = repository.NewProductPlanRepository(ptr.transactionContext)
productGroupRepository, _ = repository.NewProductGroupRepository(ptr.transactionContext)
//employeeProductRecordRepository, _ = repository.NewEmployeeProductRecordRepository(ptr.transactionContext)
)
var (
... ... @@ -60,9 +58,9 @@ func (ptr *PGProductRecordService) EmployeeProductStatics(productRecord *domain.
case ProductSection3:
break
case ProductSection4: //个人特殊处理
return ptr.personalProductStatics(nil, nil, nil, productRecord)
return ptr.personalProductStatics(nil, productRecord)
default:
return nil, nil //ptr.personalProductStatics(productRecord)
return nil, nil
}
if planId == 0 {
log.Logger.Debug(fmt.Sprintf("工段:%v product_record 编号:%v 批次为0", productRecord.WorkStation.WorkStationId, productRecord.ProductRecordId))
... ... @@ -74,9 +72,6 @@ func (ptr *PGProductRecordService) EmployeeProductStatics(productRecord *domain.
return nil, err
}
// 2.1 判断是否是支援类型 有打卡记录,员工是否是属于该工段的员工
groupMembers, groupMembersKeyFunc := FindGroupMembers(productGroupRepository, cid, oid, productRecord.WorkStation.WorkStationId)
// 集体
// 1.查询员工 -》 员工打卡记录 工位+打卡日期
// 2.打卡记录的时间区间 在生产记录上报的时间范围内
... ... @@ -91,8 +86,8 @@ func (ptr *PGProductRecordService) EmployeeProductStatics(productRecord *domain.
productRecord.ProductWorker = r.ProductWorker
// 3.查询员工产能记录 -》员工 批次+工位+批次生产日期 (有 更新产能数据、没有插入一条产能数据)
// 4.更新产能 (产能、二级品) (特殊工段处理 打料、成型)
// 个人
if _, err := ptr.personalProductStatics(productPlan, groupMembers, groupMembersKeyFunc, productRecord); err != nil {
// 个人产能统计
if _, err := ptr.personalProductStatics(productPlan, productRecord); err != nil {
return nil, err
}
}
... ... @@ -174,9 +169,8 @@ func FindGroupMembers(productGroupRepository domain.ProductGroupRepository, comp
}
// 个人生产记录统计
func (ptr *PGProductRecordService) personalProductStatics(productPlan *domain.ProductPlan, groupMembers map[string]*domain.User, groupMembersKeyFunc func(int) string, productRecord *domain.ProductRecord) (interface{}, error) {
func (ptr *PGProductRecordService) personalProductStatics(productPlan *domain.ProductPlan, productRecord *domain.ProductRecord) (interface{}, error) {
var (
//workshopRepository,_=repository.NewWorkshopRepository(ptr.transactionContext)
productPlanRepository, _ = repository.NewProductPlanRepository(ptr.transactionContext)
productGroupRepository, _ = repository.NewProductGroupRepository(ptr.transactionContext)
employeeProductRecordRepository, _ = repository.NewEmployeeProductRecordRepository(ptr.transactionContext)
... ... @@ -185,20 +179,18 @@ func (ptr *PGProductRecordService) personalProductStatics(productPlan *domain.Pr
cid = productRecord.CompanyId
oid = productRecord.OrgId
planId = productRecord.ProductRecordInfo.ProductPlanId
//productPlan *domain.ProductPlan
err error
workStationId = productRecord.WorkStation.WorkStationId
)
// 2.1 判断是否是支援类型 有打卡记录,员工是否是属于该工段的员工
if groupMembers == nil {
groupMembers, groupMembersKeyFunc = FindGroupMembers(productGroupRepository, cid, oid, productRecord.WorkStation.WorkStationId)
}
groupMembers, groupMembersKeyFunc := FindGroupMembers(productGroupRepository, cid, oid, workStationId)
if productPlan == nil {
productPlan, err = productPlanRepository.FindOne(map[string]interface{}{"productPlanId": planId})
if err != nil {
return nil, err
}
}
workStationId := productRecord.WorkStation.WorkStationId
employeeProductRecordDao, _ := dao.NewEmployeeProductRecordDao(ptr.transactionContext)
... ... @@ -231,11 +223,13 @@ func (ptr *PGProductRecordService) personalProductStatics(productPlan *domain.Pr
}
}
}
log.Logger.Debug(fmt.Sprintf("产能统计 员工:%v(%v) 当前产能:%v kg 新增产能:%v kg 原始数量:%v",
employeeProductRecord.ProductWorker.UserName, employeeProductRecord.ProductWorker.UserId,
employeeProductRecord.ProductWeigh, productRecord.ProductRecordInfo.Weigh, productRecord.ProductRecordInfo.Original))
employeeProductRecord.UpdateProductWeigh(productRecord, yesterdayOutputWeight, bestOutputWeight)
if employeeProductRecord, err = employeeProductRecordRepository.Save(employeeProductRecord); err != nil {
// TODO:异常处理
log.Logger.Error(fmt.Sprintf("生产记录:[%v] 员工:[%v] 处理异常:%v", productRecord.ProductRecordId, productRecord.ProductWorker.UserId, err.Error()))
}
return nil, nil
... ... @@ -246,7 +240,6 @@ func (ptr *PGProductRecordService) WorkshopProductStatics(productRecord *domain.
var (
workshopRepository, _ = repository.NewWorkshopRepository(ptr.transactionContext)
productPlanRepository, _ = repository.NewProductPlanRepository(ptr.transactionContext)
//productGroupRepository, _ = repository.NewProductGroupRepository(ptr.transactionContext)
workshopProductRecordRepository, _ = repository.NewWorkshopProductRecordRepository(ptr.transactionContext)
)
... ... @@ -295,6 +288,9 @@ func (ptr *PGProductRecordService) WorkshopProductStatics(productRecord *domain.
return nil, nil
}
}
log.Logger.Debug(fmt.Sprintf("产能统计 工位:%v(%v) 当前产能:%v kg 新增产能:%v kg 原始数量:%v",
workshopProductRecord.WorkStation.WorkStationId, workshopProductRecord.WorkStation.SectionName,
workshopProductRecord.ProductWeigh, productRecord.ProductRecordInfo.Weigh, productRecord.ProductRecordInfo.Original))
workshopProductRecord.UpdateProductWeigh(productRecord)
// 打料 跟 成型工段的初始产能是批次的产能
if productRecord.WorkStation.SectionName == ProductSection1 && productRecord.WorkStation.SectionName == ProductSection2 {
... ...
... ... @@ -96,10 +96,11 @@ func (ptr *PGWorkshopDataConsumeService) Consume(companyId, orgId int, record *d
if record.DeviceType == domain.DeviceTypeChuanChuanJi && plan != nil && deviceRunningData.Count > 0 {
productRecord, _ := ptr.newProductRecord(companyId, orgId, workStation, device, deviceRunningData, plan)
//if _, err = deviceRunningRecordRepository.Save(deviceRunningRecord); err != nil {
// return nil, err
//}
SendProductRecordStaticsJob(productRecord)
//SendProductRecordStaticsJob(productRecord)
productRecordService, _ := NewPGProductRecordService(ptr.transactionContext)
productRecordService.EmployeeProductStatics(productRecord)
productRecordService.WorkshopProductStatics(productRecord)
}
// 3.更新 设备每日运行记录(汇总) - redis更新 十分钟异步刷库
... ...
... ... @@ -5,6 +5,7 @@ import (
"github.com/hibiken/asynq"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/constant"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/domain"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/infrastructure/domainService"
"gitlab.fjmaimaimai.com/allied-creation/allied-creation-manufacture/pkg/log"
"os"
"os/signal"
... ... @@ -20,12 +21,15 @@ func Run() {
srv := asynq.NewServer(
asynq.RedisClientOpt{Addr: constant.REDIS_ADDRESS},
asynq.Config{
Concurrency: 2,
Concurrency: 4,
Queues: map[string]int{
//"critical": 1,
"default": 1,
"device": 1,
domainService.FormatQueue(domainService.QueueDevice): 1,
domainService.FormatQueue(domainService.QueueProduct): 1,
domainService.FormatQueue(domainService.QueueDefault): 1,
},
StrictPriority: true,
},
)
... ...
... ... @@ -17,8 +17,8 @@ func HandlerProductRecordStatics(c context.Context, t *asynq.Task) error {
if err := json.Unmarshal(t.Payload(), cmd); err != nil {
return err
}
log.Logger.Debug(fmt.Sprintf("【生产记录统计】 消费 生产记录ID:%v 类型:%v 工段:%v(%v) 重量:%v", cmd.ProductRecordId, cmd.ProductRecordType,
cmd.WorkStation.SectionName, cmd.WorkStation.WorkStationId, cmd.ProductRecordInfo.Weigh))
log.Logger.Debug(fmt.Sprintf("【生产记录统计】 消费 生产记录ID:%v 类型:%v 工段:%v(%v) 重量:%v 记录时间:%v", cmd.ProductRecordId, cmd.ProductRecordType,
cmd.WorkStation.SectionName, cmd.WorkStation.WorkStationId, cmd.ProductRecordInfo.Weigh, cmd.CreatedAt))
_, err := productPlanService.ProductRecordStatics(cmd)
if err != nil {
log.Logger.Error(err.Error())
... ...