run_task.go 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371
  1. package application
  2. import (
  3. "dcl_dispatch_server/src/package/domain/project"
  4. "dcl_dispatch_server/src/package/domain/task_cache"
  5. "dcl_dispatch_server/src/package/infra"
  6. "dcl_dispatch_server/src/package/infra/global"
  7. "dcl_dispatch_server/src/package/infra/redis"
  8. "dcl_dispatch_server/src/package/util"
  9. "fmt"
  10. "github.com/aliyun/aliyun-oss-go-sdk/oss"
  11. "path/filepath"
  12. "strconv"
  13. "strings"
  14. "time"
  15. "github.com/confluentinc/confluent-kafka-go/kafka"
  16. )
  17. /*
  18. 负责处理用户等待队列中的任务
  19. 负责运行集群等待队列中的任务
  20. */
  21. // 判断用户等待队列中的任务是否可以加入到集群等待队列
  22. func RunWaitingUser() {
  23. infra.GlobalLogger.Infof("启动【用户等待队列】监控进程。")
  24. for {
  25. time.Sleep(2 * time.Second)
  26. global.RunTaskMutex.Lock()
  27. // 获取Redis列表中的值
  28. taskCacheJsons, err := infra.GlobalRedisClient.LRange(global.KeyTaskQueueWaitingUser, 0, -1).Result()
  29. if err != nil {
  30. infra.GlobalLogger.Errorf("遍历用户等待队列 %v 失败,错误信息为: %v", global.KeyTaskQueueWaitingUser, err)
  31. continue
  32. }
  33. for _, taskCacheJson := range taskCacheJsons {
  34. taskCache, err := task_cache.JsonToTaskCache(taskCacheJson)
  35. if err != nil {
  36. infra.GlobalLogger.Error(err)
  37. continue
  38. }
  39. userId := taskCache.UserId
  40. userParallelism := taskCache.UserParallelism
  41. algorithmObjectKey := taskCache.AlgorithmObjectKey
  42. equipmentType := taskCache.EquipmentType
  43. task := taskCache.Task
  44. // 1 判断用户并行度是否有剩余,有剩余则加入集群等待队列,并从用户等待队列中拿出,没有剩余则不需要改动
  45. if redis.CanRunUser(userId, userParallelism) { // 可以运行
  46. err = redis.AddWaitingCluster(&project.Project{
  47. UserId: userId,
  48. Parallelism: userParallelism,
  49. AlgorithmObjectKey: algorithmObjectKey,
  50. EquipmentType: equipmentType,
  51. }, task)
  52. if err != nil {
  53. infra.GlobalLogger.Error(err)
  54. continue
  55. }
  56. err = redis.DeleteWaitingUser(task.Info.TaskId)
  57. if err != nil {
  58. infra.GlobalLogger.Error(err)
  59. continue
  60. }
  61. }
  62. }
  63. global.RunTaskMutex.Unlock()
  64. }
  65. }
  66. // 集群等待队列中的任务判断是否可以加入集群运行队列
  67. func RunWaitingCluster() {
  68. infra.GlobalLogger.Infof("启动【集群等待队列】监控进程。")
  69. for {
  70. time.Sleep(2 * time.Second)
  71. global.GpuNodeListMutex.Lock()
  72. // 1 判断用户并行度是否有剩余,有剩余则从集群等待队列取出第一个加入集群运行队列,并运行pod,没有剩余则不需要改动
  73. can, gpuNode, err := redis.CanRunCluster()
  74. if err != nil {
  75. infra.GlobalLogger.Error(err)
  76. global.GpuNodeListMutex.Unlock()
  77. continue
  78. }
  79. firstTaskCache := task_cache.TaskCache{}
  80. algorithmTarName := ""
  81. algorithmTarPath := ""
  82. algorithmImageName := ""
  83. algorithmImageNameWithVersion := ""
  84. algorithmExist := false
  85. // 确定用哪个oss
  86. var tempOss *oss.Bucket
  87. if can {
  88. // 判断是否有待运行的任务
  89. waitingClusterNumber, _ := infra.GlobalRedisClient.LLen(global.KeyTaskQueueWaitingCluster).Result()
  90. if waitingClusterNumber == 0 {
  91. global.GpuNodeListMutex.Unlock()
  92. continue
  93. } else {
  94. infra.GlobalLogger.Infof("集群存在 %v 个等待运行的任务。", waitingClusterNumber)
  95. }
  96. // 取出并移出,20241017 防止报错后陷入死循环
  97. {
  98. firstTaskCacheJson, err := infra.GlobalRedisClient.LIndex(global.KeyTaskQueueWaitingCluster, 0).Result()
  99. if err != nil {
  100. infra.GlobalLogger.Error("取出集群等待队列中的头元素报错,错误信息为:", err)
  101. global.GpuNodeListMutex.Unlock()
  102. continue
  103. }
  104. firstTaskCache, err = task_cache.JsonToTaskCache(firstTaskCacheJson)
  105. if err != nil {
  106. infra.GlobalLogger.Error(err)
  107. global.GpuNodeListMutex.Unlock()
  108. continue
  109. }
  110. // --------------- 从等待队列中移除
  111. if _, err = infra.GlobalRedisClient.LPop(global.KeyTaskQueueWaitingCluster).Result(); err != nil {
  112. infra.GlobalLogger.Error(err)
  113. continue
  114. }
  115. }
  116. if firstTaskCache.Env == global.EnvPji {
  117. tempOss = infra.GlobalOssBucketPji
  118. } else {
  119. tempOss = infra.GlobalOssBucketCicv
  120. }
  121. // --------------- 下载算法 ---------------
  122. {
  123. infra.GlobalLogger.Infof("开始下载算法 %v。", firstTaskCache.AlgorithmObjectKey)
  124. algorithmTarName = filepath.Base(firstTaskCache.AlgorithmObjectKey)
  125. algorithmTarPath = infra.ApplicationYaml.K8s.AlgorithmTarTempDir + util.NewShortUUID() + "/" + algorithmTarName
  126. _ = util.CreateParentDir(algorithmTarPath)
  127. algorithmImageName = infra.ApplicationYaml.K8s.RegistryUri + "/cicvdcl_" + util.MD5HashShort(algorithmTarName)
  128. algorithmImageNameWithVersion = algorithmImageName + ":latest"
  129. algorithmExist = util.ImageExists(infra.GlobalDockerClient, algorithmImageName)
  130. if !algorithmExist {
  131. err = tempOss.GetObjectToFile(firstTaskCache.AlgorithmObjectKey, algorithmTarPath)
  132. if err != nil {
  133. infra.GlobalLogger.Error("下载oss上的算法镜像 "+firstTaskCache.AlgorithmObjectKey+" 失败,错误信息为:", err)
  134. time.Sleep(time.Duration(2) * time.Second)
  135. global.GpuNodeListMutex.Unlock()
  136. continue
  137. }
  138. infra.GlobalLogger.Infof("下载算法 %v 成功。", firstTaskCache.AlgorithmObjectKey)
  139. } else {
  140. infra.GlobalLogger.Infof("算法 %v 已存在。", algorithmImageName)
  141. }
  142. }
  143. } else {
  144. infra.GlobalLogger.Infof("集群没有剩余并行度。")
  145. global.GpuNodeListMutex.Unlock()
  146. continue
  147. }
  148. global.GpuNodeListMutex.Unlock()
  149. // 获取项目ID
  150. projectId := firstTaskCache.Task.Info.ProjectId
  151. offsetKey := "offset:" + projectId
  152. offset := 0
  153. // 根据项目ID获取偏移量
  154. val, err := infra.GlobalRedisClient.Get(offsetKey).Result()
  155. if err != nil {
  156. //infra.GlobalLogger.Infof("偏移量键 %v 不存在,初始化设置为 0。", offsetKey)
  157. err = infra.GlobalRedisClient.Set(offsetKey, 0, 0).Err()
  158. if err != nil {
  159. infra.GlobalLogger.Infof("偏移量键值对 %v 初始化失败,错误信息为: %v", offsetKey, err)
  160. continue
  161. }
  162. } else {
  163. offset, err = strconv.Atoi(val)
  164. if err != nil {
  165. infra.GlobalLogger.Infof("字符串 %v 转整数失败,错误信息为: %v", val, err)
  166. continue
  167. }
  168. infra.GlobalLogger.Infof("当前任务使用偏移量【%v】", offset)
  169. }
  170. // 取出偏移量后将缓存中的加一,给下个任务使用。
  171. _, err = infra.GlobalRedisClient.Incr(offsetKey).Result()
  172. if err != nil {
  173. infra.GlobalLogger.Infof("偏移量 %v 加一失败,错误信息为: %v", offsetKey, err)
  174. continue
  175. }
  176. //infra.GlobalLogger.Infof("偏移量【%v】加一给下个任务使用。", offsetKey)
  177. // ------- 解析 xosc ,从xosc中解析起终点位置&修改 xodr 和 osgb ---------------
  178. xoscOssPath := firstTaskCache.Task.Scenario.ScenarioOsc
  179. tempDir := infra.ApplicationYaml.TempDir
  180. util.CreateDir(tempDir)
  181. xoscLocalPath := tempDir + util.NewShortUUID() + ".xosc"
  182. err = tempOss.GetObjectToFile(firstTaskCache.Task.Scenario.ScenarioOsc, xoscLocalPath)
  183. if err != nil {
  184. infra.GlobalLogger.Errorf("下载xosc文件【%v】失败,错误信息为:%v", xoscOssPath, err)
  185. continue
  186. }
  187. startPositionX, startPositionY, startPositionZ, startPositionH, _, _, _, _, xodrPath, osgbPath := util.ParseXosc(xoscLocalPath)
  188. firstTaskCache.Task.Scenario.ScenarioOdr = xodrPath
  189. firstTaskCache.Task.Scenario.ScenarioOsgb = osgbPath
  190. destinationCsvOssKey := getCsvByXosc(xoscOssPath)
  191. var endPoints []util.DataPoint
  192. if exist, _ := tempOss.IsObjectExist(destinationCsvOssKey); exist {
  193. destinationCsvLocalPath := tempDir + util.NewShortUUID() + ".xosc"
  194. err = tempOss.GetObjectToFile(destinationCsvOssKey, destinationCsvLocalPath)
  195. endPoints, _ = util.ParseCsv(destinationCsvLocalPath)
  196. infra.GlobalLogger.Infof("终点文件【%v】的解析结果为:【%v】", destinationCsvLocalPath, endPoints)
  197. } else {
  198. endPositionX, endPositionY, endPositionZ, endPositionH := global.DefaultEndPositionX, global.DefaultEndPositionY, global.DefaultEndPositionZ, global.DefaultEndPositionH
  199. tempEndPoint := util.DataPoint{X: endPositionX, Y: endPositionY, Z: endPositionZ, H: endPositionH}
  200. endPoints = append(endPoints, tempEndPoint)
  201. infra.GlobalLogger.Infof("终点文件【%v】不存在,使用默认终点。", destinationCsvOssKey)
  202. }
  203. var endPosition string // export END_POSITION="x1,y1,z1,h1;x2,y2,z2,h2;x3,y3,z3,h3;x4,y4,z4,h4"
  204. // 遍历生成字符串
  205. var endPositionBuilder strings.Builder
  206. for _, endPoint := range endPoints {
  207. if endPositionBuilder.Len() > 0 {
  208. endPositionBuilder.WriteByte(';')
  209. }
  210. endPositionBuilder.WriteString(fmt.Sprintf("%s,%s,%s,%s", endPoint.X, endPoint.Y, endPoint.Z, endPoint.H))
  211. }
  212. endPosition = endPositionBuilder.String()
  213. // ------- 修改OGT传感器显示目标物框 -------
  214. for i := range firstTaskCache.Task.Vehicle.Sensors.OGT {
  215. firstTaskCache.Task.Vehicle.Sensors.OGT[i].SensorDisplay = global.DefaultOgtSensorDisplay
  216. }
  217. // --------------- 发送 kafka 消息(获取偏移量和分区) ---------------
  218. // 获取任务消息转json
  219. firstTaskCache.Task.Vehicle.Sensors.Camera = global.DefaultCameras
  220. taskJson, err := project.TaskToJson(firstTaskCache.Task)
  221. if err != nil {
  222. infra.GlobalLogger.Error(err)
  223. continue
  224. }
  225. // ------- 发送 -------
  226. topic := projectId
  227. // 创建一个Message,并指定分区为0
  228. msg := &kafka.Message{
  229. TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: infra.ApplicationYaml.Kafka.Partition, Offset: kafka.Offset(offset)},
  230. Value: []byte(taskJson),
  231. }
  232. // 发送消息,并处理结果
  233. err = infra.GlobalKafkaProducer.Produce(msg, nil)
  234. if err != nil {
  235. infra.GlobalLogger.Infof("发送任务消息 %v 失败,错误信息为: %v", msg, err)
  236. continue
  237. }
  238. infra.GlobalLogger.Infof("发送任务消息成功,话题为【%v】,偏移量为【%v】,消息为【%v】", topic, offset, taskJson)
  239. // 如果新算法需要导入
  240. if !algorithmExist {
  241. // 导入算法
  242. infra.GlobalLogger.Infof("导入算法文件【%v】到docker镜像【%v】。", algorithmTarPath, algorithmImageNameWithVersion)
  243. _, s, err := util.Execute("docker", "import", algorithmTarPath, algorithmImageNameWithVersion)
  244. infra.GlobalLogger.Infof("推送算法镜像【%v】。", algorithmImageNameWithVersion)
  245. _, s, err = util.Execute("docker", "push", algorithmImageNameWithVersion)
  246. if err != nil {
  247. infra.GlobalLogger.Errorf("导入算法镜像 %v 为 %v 失败,执行结果为:%v,错误信息为:%v", algorithmTarPath, algorithmImageNameWithVersion, s, err)
  248. time.Sleep(time.Duration(2) * time.Second)
  249. continue
  250. }
  251. infra.GlobalLogger.Infof("导入算法镜像 %v 为 %v 成功,执行结果为:%v", algorithmTarPath, algorithmImageNameWithVersion, s)
  252. err = util.RemoveFile(algorithmTarPath)
  253. if err != nil {
  254. infra.GlobalLogger.Errorf("删除算法镜像文件 %v 失败,错误信息为:%v", algorithmTarPath, err)
  255. }
  256. }
  257. // --------------- 启动 k8s pod ---------------
  258. podName := "project-" + projectId + "-" + util.NewShortUUID()
  259. namespaceName := infra.ApplicationYaml.K8s.NamespaceName
  260. nodeName := gpuNode.Hostname
  261. restParallelism := gpuNode.Parallelism
  262. vtdContainer := "vtd-" + projectId
  263. algorithmContainer := "algorithm-" + projectId
  264. vtdImage := infra.ApplicationYaml.K8s.VtdImage
  265. // 2 生成模板文件名称
  266. podYaml := nodeName + "#" + podName + ".yaml"
  267. // 3 模板yaml存储路径
  268. yamlPath := infra.ApplicationYaml.K8s.PodYamlDir + podYaml
  269. // 4 模板yaml备份路径
  270. yamlPathBak := infra.ApplicationYaml.K8s.PodYamlDir + "bak/" + podYaml
  271. fmt.Println(yamlPath, yamlPathBak)
  272. // 5 不同设备类型用不同的模板文件
  273. podString := ""
  274. if firstTaskCache.EquipmentType == global.EquipmentTypeKinglong || firstTaskCache.EquipmentType == global.EquipmentTypePjisuv { // 多功能车仿真
  275. if podString, err = util.ReadFile(infra.ApplicationYaml.K8s.VtdPodTemplateYamlPjisuv); err != nil {
  276. infra.GlobalLogger.Error(err)
  277. continue
  278. }
  279. podString = strings.Replace(podString, "vtd-command", infra.ApplicationYaml.K8s.VtdCommandPjisuv, -1)
  280. } else {
  281. if podString, err = util.ReadFile(infra.ApplicationYaml.K8s.VtdPodTemplateYamlPjibot); err != nil {
  282. infra.GlobalLogger.Error(err)
  283. continue
  284. }
  285. podString = strings.Replace(podString, "vtd-command", infra.ApplicationYaml.K8s.VtdCommandPjibot, -1)
  286. }
  287. if firstTaskCache.Env == global.EnvPji {
  288. podString = strings.Replace(podString, "oss-type", infra.ApplicationYaml.OssPji.Type, -1)
  289. podString = strings.Replace(podString, "oss-ip", infra.ApplicationYaml.OssPji.Endpoint, -1) // 不带http://前缀
  290. podString = strings.Replace(podString, "oss-access-key", infra.ApplicationYaml.OssPji.AccessKeyId, -1)
  291. podString = strings.Replace(podString, "oss-secret-key", infra.ApplicationYaml.OssPji.AccessKeySecret, -1)
  292. podString = strings.Replace(podString, "oss-bucket", infra.ApplicationYaml.OssPji.BucketName, -1)
  293. } else {
  294. podString = strings.Replace(podString, "oss-type", infra.ApplicationYaml.OssCicv.Type, -1)
  295. podString = strings.Replace(podString, "oss-ip", infra.ApplicationYaml.OssCicv.Endpoint, -1) // 不带http://前缀
  296. podString = strings.Replace(podString, "oss-access-key", infra.ApplicationYaml.OssCicv.AccessKeyId, -1)
  297. podString = strings.Replace(podString, "oss-secret-key", infra.ApplicationYaml.OssCicv.AccessKeySecret, -1)
  298. podString = strings.Replace(podString, "oss-bucket", infra.ApplicationYaml.OssCicv.BucketName, -1)
  299. }
  300. podString = strings.Replace(podString, "pod-name", podName, -1)
  301. podString = strings.Replace(podString, "namespace-name", namespaceName, -1)
  302. podString = strings.Replace(podString, "node-name", nodeName, -1)
  303. podString = strings.Replace(podString, "algorithm-image", algorithmImageNameWithVersion, -1)
  304. podString = strings.Replace(podString, "vtd-container", vtdContainer, -1)
  305. podString = strings.Replace(podString, "vtd-image", vtdImage, -1)
  306. podString = strings.Replace(podString, "platform-ip", infra.ApplicationYaml.Web.IpPrivate+":"+infra.ApplicationYaml.Web.Port, -1)
  307. podString = strings.Replace(podString, "simulation-cloud-ip", infra.ApplicationYaml.Web.IpPrivate+":"+infra.ApplicationYaml.Web.Port, -1)
  308. podString = strings.Replace(podString, "platform-type", "\""+infra.ApplicationYaml.K8s.PlatformType+"\"", -1)
  309. podString = strings.Replace(podString, "kafka-ip", infra.ApplicationYaml.Kafka.Broker, -1)
  310. podString = strings.Replace(podString, "kafka-topic", projectId, -1)
  311. podString = strings.Replace(podString, "kafka-partition", "\""+util.ToString(infra.ApplicationYaml.Kafka.Partition)+"\"", -1)
  312. podString = strings.Replace(podString, "kafka-offset", "\""+util.ToString(offset)+"\"", -1)
  313. podString = strings.Replace(podString, "cpu-order", "\""+util.ToString(restParallelism-1)+"\"", -1) // cpu编号是剩余并行度-1
  314. podString = strings.Replace(podString, "algorithm-container", algorithmContainer, -1)
  315. podString = strings.Replace(podString, "start-position-x", "\""+startPositionX+"\"", -1)
  316. podString = strings.Replace(podString, "start-position-y", "\""+startPositionY+"\"", -1)
  317. podString = strings.Replace(podString, "start-position-z", "\""+startPositionZ+"\"", -1)
  318. podString = strings.Replace(podString, "start-position-h", "\""+startPositionH+"\"", -1)
  319. podString = strings.Replace(podString, "end-position", "\""+endPosition+"\"", -1)
  320. // --------------- 保存成文件 ----------------
  321. err = util.WriteFile(podString, yamlPath)
  322. err = util.WriteFile(podString, yamlPathBak)
  323. if err != nil {
  324. infra.GlobalLogger.Error("保存yaml字符串失败,错误信息为", err)
  325. continue
  326. }
  327. infra.GlobalLogger.Infof("保存yaml文件到执行路径【%v】和备份路径【%v】", yamlPath, yamlPathBak)
  328. // --------------- 启动 pod
  329. _, sr, err := util.Execute("kubectl", "apply", "-f", yamlPath)
  330. if err != nil {
  331. infra.GlobalLogger.Errorf("启动pod失败,执行结果为 %v,错误信息为 %v", sr, err)
  332. continue
  333. }
  334. infra.GlobalLogger.Errorf("启动pod成功,执行结果为 %v。", sr)
  335. // 收尾
  336. {
  337. // --------------- 添加到运行队列
  338. err = redis.AddRunningCluster(firstTaskCache, gpuNode.Hostname)
  339. if err != nil {
  340. infra.GlobalLogger.Error(err)
  341. global.GpuNodeListMutex.Unlock()
  342. continue
  343. }
  344. // --------------- 删除镜像文件
  345. _ = util.RemoveFile(algorithmTarPath)
  346. }
  347. }
  348. }
  349. func getCsvByXosc(xoscOssKey string) string {
  350. lastIndex := strings.LastIndex(xoscOssKey, "/")
  351. dirPath := xoscOssKey[:lastIndex+1] // 包括最后一个 /
  352. destinationCsvOssKey := dirPath + global.DestinationCsvName
  353. return destinationCsvOssKey
  354. }