| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427 |
- using FlowAlert.Model;
- using log4net;
- using Newtonsoft.Json;
- using Newtonsoft.Json.Linq;
- using Quartz;
- using RDIFramework.Utilities;
- using System;
- using System.Collections.Generic;
- using System.Configuration;
- using System.Data;
- using System.Data.SqlClient;
- using System.IO;
- using System.Linq;
- using System.Net;
- using System.Security.Cryptography;
- using System.Text;
- using System.Web.Security;
- namespace TimedUpload.QuartzJobs
- {
- /// <summary>
- /// DisallowConcurrentExecution:同一调度器内上一次执行未结束时,不允许再次并发执行
- /// </summary>
- [DisallowConcurrentExecution]
- public class DataUploadJob : IJob
- {
- private readonly ILog log = LogManager.GetLogger(typeof(DataUploadJob));
- static IDbProvider dbHelper
- {
- get
- {
- var DbDefine = DbFactoryProvider.GetProvider(CurrentDbType.SqlServer, Config.GetValue("DbConn"));
- return DbDefine;
- }
- }
- public void Execute(IJobExecutionContext context)
- {
- try
- {
- String uploadDays = Config.GetValue("UploadDays");
- string[] daysArr = uploadDays.Split(',');
- string day = DateTime.Now.Day.ToString();
- if (daysArr.Contains(day))
- {
- UploadDataInfo();
- }
- else
- {
- log.Info("不在上传数据日期内");
- }
- }
- catch (Exception ex)
- {
- log.Error("上传数据错误:" + ex.Message);
- }
- }
- /// <summary>
- /// 上传数据
- /// </summary>
- public void UploadDataInfo()
- {
- String result = "";
-
- String orgid = Config.GetValue("orgid");
- String key = Config.GetValue("key");
- String value;
- try
- {
- StringBuilder stb = new StringBuilder();
- String sql = " select UserNo,UserName,totalAddr,ElecAddress,PhoneNumber,NowRadingDT,CONVERT(int, NowReading) 表读数 from dbo.V_WaterMetersAll where companyid='fe5b322f-a343-43bd-b622-171692ee6197' and UserNo<>'' ";
-
- DataTable dt = dbHelper.Fill(sql);
- if (dt != null && dt.Rows.Count > 0)
- {
- int num = 0;
- int i = 0;
- int skip = 0;
- string readingMonth = DateTime.Now.ToString("yyyyMM");
- // 防重:一次性查出本月所有已成功推送的读数版本(水表/用户 + 读数 + 抄表时间)。
- // 同一读数版本已成功推过才跳过;月内新读数是新版本,会继续推送;失败不进版本集合,下次补传。
- // 推送记录表为追加流水:每次推送(成功/失败、月内第几次数)都新增一条,不更新历史记录。
- HashSet<PushedVersion> pushedVersions = GetPushedVersions(readingMonth);
- log.Debug("~~~~~~~~~~~开始上传数据~~~~~~~~~~~~~");
- foreach (DataRow row in dt.Rows)
- {
- try
- {
- i += 1;
- string userNo = row["UserNo"] == DBNull.Value ? string.Empty : row["UserNo"].ToString();
- string meterNo = row["ElecAddress"] == DBNull.Value ? string.Empty : row["ElecAddress"].ToString();
- int nowReadingVal = row["表读数"] == DBNull.Value ? 0 : Convert.ToInt32(row["表读数"]);
- DateTime nowRadingDTVal = row["NowRadingDT"] == DBNull.Value ? DateTime.Now : Convert.ToDateTime(row["NowRadingDT"]);
- string userName = row["UserName"] == DBNull.Value ? string.Empty : row["UserName"].ToString();
- string totalAddr = row["totalAddr"] == DBNull.Value ? string.Empty : row["totalAddr"].ToString();
- string phoneNumber = row["PhoneNumber"] == DBNull.Value ? string.Empty : row["PhoneNumber"].ToString();
- // 版本判重:该读数+抄表时间版本本月已成功推过则跳过;读数更新产生新版本,继续推送并另行留痕
- PushedVersion currentVersion = new PushedVersion(BuildDedupKey(meterNo, userNo), nowReadingVal, nowRadingDTVal);
- if (pushedVersions != null && pushedVersions.Contains(currentVersion))
- {
- skip += 1;
- log.Debug("本月(" + readingMonth + ")该读数版本已成功上报,跳过--------》" + meterNo
- + ", 读数:" + nowReadingVal + "@" + nowRadingDTVal.ToString("yyyy-MM-dd"));
- continue;
- }
- StringBuilder dataStb = new StringBuilder();
- dataStb.Clear();
- dataStb.Append("[{");
- dataStb.Append("\"CBNY\":\"" + readingMonth + "\",");
- dataStb.Append("\"DATA\":[{");
- dataStb.Append("\"SYBS\":\"0\",");
- dataStb.Append("\"CBRQ\":\"" + nowRadingDTVal.ToString("yyyy-MM-dd") + "\",");
- dataStb.Append("\"SBBH\":\"" + row["ElecAddress"] + "\",");
- dataStb.Append("\"CUSERID\":\"" + row["UserNo"] + "\",");
- dataStb.Append("\"CBZT\":\"正常\",");
- dataStb.Append("\"BYBS\":" + row["表读数"] + ",");
- dataStb.Append("\"CCOPY\":" + 1 + ",");
- dataStb.Append("\"JJSL\":" + 0.0 + ",");
- dataStb.Append("\"REMARK\":\"\",");
- dataStb.Append("}],");
- dataStb.Append("\"JYDM\":\"024\",");
- dataStb.Append("\"CJDM\":\"WWKJ\"");
- dataStb.Append("}]");
- log.Debug("上传数据---------》" + row["ElecAddress"] + "," + dataStb.ToString());
- result = postSend(Config.GetValue("UploadDataUrl"), dataStb.ToString());
- dataStb.Clear();
- log.Debug("上传结果---------》" + row["ElecAddress"] + "," + result);
-
- // 解析上传结果
- bool success = false;
- string resultMsg = string.Empty;
- if (!string.IsNullOrEmpty(result))
- {
- // var dicRes = JsonConvert.DeserializeObject<Result>(result);
- var dicRes = JArray.Parse(result);
- if (dicRes != null && dicRes.Count > 0)
- {
- success = dicRes.First["SUCCESS"].ToString() == "True";
- resultMsg = dicRes.First["MSG"] != null ? dicRes.First["MSG"].ToString() : string.Empty;
- }
- }
- else
- {
- resultMsg = "请求失败,未返回结果";
- }
-
- if (success)
- {
- log.Debug(row["ElecAddress"] + "上传成功");
- num += 1;
- }
- else
- {
- log.Debug(row["ElecAddress"] + "上传失败");
- }
-
- // 追加写入推送记录表:每次推送都新增一条流水(成功/失败均留痕,月内多个读数版本各有各的记录)
- InsertPushRecord(
- userNo,
- userName,
- totalAddr,
- meterNo,
- phoneNumber,
- nowRadingDTVal,
- readingMonth,
- nowReadingVal,
- success ? "成功" : "失败",
- resultMsg);
- if (success)
- {
- // 本次成功的版本纳入内存集合,避免同一次执行内重复推送同一版本
- if (pushedVersions != null)
- {
- pushedVersions.Add(currentVersion);
- }
- }
-
- if (i == dt.Rows.Count)
- {
- log.Debug("总共" + dt.Rows.Count + "条, 成功上传" + num + "条, 跳过已上报" + skip + "条, 上传失败:" + (dt.Rows.Count - num - skip));
- }
- }
- catch (Exception ex)
- {
- log.Debug("上传数据失败:" + row["ElecAddress"] + "," + ex.StackTrace );
- }
-
- }
- }
- }
- catch (Exception ex)
- {
- log.Debug("上传数据失败:" + ex.ToString());
- }
- }
-
- #region 请求接口
- /// <summary>
- ///
- /// </summary>
- /// <param name="url"></param>
- /// <param name="param"></param>
- /// <param name="type"></param>
- /// <param name="webHeader">请求头携带请求权限</param>
- /// <returns></returns>
- public string postSend(string url, string param, String type = "POST")
- {
- string strResult = "";
- Encoding myEncode = Encoding.GetEncoding("UTF-8");
- HttpWebRequest req = (HttpWebRequest)HttpWebRequest.Create(url);
- req.Method = type;
- req.ContentType = "application/json; charset=utf-8";
-
-
- if (param != null)
- {
- byte[] postBytes = Encoding.UTF8.GetBytes(param);
- req.ContentLength = postBytes.Length;
- using (Stream reqStream = req.GetRequestStream())
- {
- reqStream.Write(postBytes, 0, postBytes.Length);
- }
- }
- try
- {
- using (WebResponse res = req.GetResponse())
- {
- using (StreamReader sr = new StreamReader(res.GetResponseStream(), myEncode))
- {
- strResult = sr.ReadToEnd();
- return strResult;
- }
- }
- }
- catch (WebException ex)
- {
- log.Error("Post数据出错:" + ex.Message);
- return "";
- }
- }
- #endregion
- #region 写入推送记录表
- /// <summary>
- /// 构造防重键:优先用水表编号(ElecAddress,对应报文SBBH),水表编号为空时回退到用户编号(UserNo)
- /// </summary>
- private static string BuildDedupKey(string meterNo, string userNo)
- {
- return string.IsNullOrEmpty(meterNo) ? "U:" + (userNo ?? string.Empty) : "M:" + meterNo;
- }
- /// <summary>
- /// 已成功推送的读数版本:防重键(水表号/用户号) + 读数 + 抄表时间 三者完全相同才算同一版本。
- /// 月内读数更新(读数或抄表时间变化)即为新版本,需要再次推送并单独留痕。
- /// </summary>
- private class PushedVersion : IEquatable<PushedVersion>
- {
- public string Key { get; private set; }
- public int NowReading { get; private set; }
- public DateTime NowRadingDT { get; private set; }
- public PushedVersion(string key, int nowReading, DateTime nowRadingDT)
- {
- Key = key ?? string.Empty;
- NowReading = nowReading;
- NowRadingDT = nowRadingDT;
- }
- public bool Equals(PushedVersion other)
- {
- if (other == null) return false;
- return string.Equals(Key, other.Key, StringComparison.OrdinalIgnoreCase)
- && NowReading == other.NowReading
- && NowRadingDT == other.NowRadingDT;
- }
- public override bool Equals(object obj)
- {
- return Equals(obj as PushedVersion);
- }
- public override int GetHashCode()
- {
- unchecked
- {
- int hash = 17;
- hash = hash * 23 + StringComparer.OrdinalIgnoreCase.GetHashCode(Key);
- hash = hash * 23 + NowReading.GetHashCode();
- hash = hash * 23 + NowRadingDT.GetHashCode();
- return hash;
- }
- }
- }
- /// <summary>
- /// 查询指定抄表月份所有已成功推送的读数版本,用于上报前按版本防重。
- /// 失败记录不包含在内,重跑时失败的读数版本会重新补传。
- /// 查询本身失败时返回 null(fail-open:宁可重复也不漏传),由调用方决定不做跳过。
- /// </summary>
- private HashSet<PushedVersion> GetPushedVersions(string readingMonth)
- {
- string connStr = Config.GetValue("DbConn");
- const string sql = @"SELECT ElecAddress, UserNo, NowReading, NowRadingDT FROM dbo.PushRecord
- WHERE ReadingMonth = @ReadingMonth AND UploadResult = N'成功'";
- HashSet<PushedVersion> versions = new HashSet<PushedVersion>();
- try
- {
- using (SqlConnection conn = new SqlConnection(connStr))
- {
- conn.Open();
- using (SqlCommand cmd = new SqlCommand(sql, conn))
- {
- cmd.Parameters.Add(new SqlParameter("@ReadingMonth", readingMonth));
- using (SqlDataReader reader = cmd.ExecuteReader())
- {
- while (reader.Read())
- {
- string meterNo = reader.IsDBNull(0) ? string.Empty : reader["ElecAddress"].ToString();
- string userNo = reader.IsDBNull(1) ? string.Empty : reader["UserNo"].ToString();
- int readingVal = reader.IsDBNull(2) ? 0 : Convert.ToInt32(reader["NowReading"]);
- DateTime readingDT = reader.IsDBNull(3) ? DateTime.MinValue : Convert.ToDateTime(reader["NowRadingDT"]);
- versions.Add(new PushedVersion(BuildDedupKey(meterNo, userNo), readingVal, readingDT));
- }
- }
- }
- }
- return versions;
- }
- catch (Exception ex)
- {
- log.Error("查询本月已上报读数版本失败,本次不做防重跳过(全部照常上报):" + ex.Message);
- return null;
- }
- }
- /// <summary>
- /// 上传后写入推送记录表 dbo.PushRecord(字段对应“每月推送模版表”,列名与数据库视图 V_WaterMetersAll 字段名一致)
- /// 使用参数化查询避免特殊字符造成 SQL 注入/语法错误
- /// </summary>
- private void InsertPushRecord(string userNo, string userName, string totalAddr, string meterNo,
- string phoneNumber, DateTime nowRadingDT, string readingMonth, int nowReading,
- string uploadResult, string resultMsg)
- {
- string connStr = Config.GetValue("DbConn");
- // 列名与数据库字段名一致:UserNo/UserName/totalAddr/ElecAddress/PhoneNumber/NowRadingDT/NowReading
- const string insertSql = @"insert into dbo.PushRecord
- (UserNo, UserName, totalAddr, ElecAddress, PhoneNumber, NowRadingDT, ReadingMonth, NowReading, UploadResult, ResultMsg, UploadTime, Remark)
- values
- (@UserNo, @UserName, @totalAddr, @ElecAddress, @PhoneNumber, @NowRadingDT, @ReadingMonth, @NowReading, @UploadResult, @ResultMsg, @UploadTime, @Remark)";
- try
- {
- using (SqlConnection conn = new SqlConnection(connStr))
- {
- conn.Open();
- using (SqlCommand cmd = new SqlCommand(insertSql, conn))
- {
- cmd.Parameters.Add(new SqlParameter("@UserNo", string.IsNullOrEmpty(userNo) ? (object)DBNull.Value : userNo));
- cmd.Parameters.Add(new SqlParameter("@UserName", string.IsNullOrEmpty(userName) ? (object)DBNull.Value : userName));
- cmd.Parameters.Add(new SqlParameter("@totalAddr", string.IsNullOrEmpty(totalAddr) ? (object)DBNull.Value : totalAddr));
- cmd.Parameters.Add(new SqlParameter("@ElecAddress", string.IsNullOrEmpty(meterNo) ? (object)DBNull.Value : meterNo));
- cmd.Parameters.Add(new SqlParameter("@PhoneNumber", string.IsNullOrEmpty(phoneNumber) ? (object)DBNull.Value : phoneNumber));
- cmd.Parameters.Add(new SqlParameter("@NowRadingDT", nowRadingDT));
- cmd.Parameters.Add(new SqlParameter("@ReadingMonth", string.IsNullOrEmpty(readingMonth) ? (object)DBNull.Value : readingMonth));
- cmd.Parameters.Add(new SqlParameter("@NowReading", nowReading));
- cmd.Parameters.Add(new SqlParameter("@UploadResult", string.IsNullOrEmpty(uploadResult) ? (object)DBNull.Value : uploadResult));
- cmd.Parameters.Add(new SqlParameter("@ResultMsg", string.IsNullOrEmpty(resultMsg) ? (object)DBNull.Value : resultMsg));
- cmd.Parameters.Add(new SqlParameter("@UploadTime", DateTime.Now));
- cmd.Parameters.Add(new SqlParameter("@Remark", (object)DBNull.Value));
- cmd.ExecuteNonQuery();
- }
- }
- }
- catch (Exception ex)
- {
- log.Error("写入推送记录表失败:" + meterNo + "," + ex.Message);
- }
- }
- #endregion
- public string md5Encript(String str) {
- MD5 md5 = new MD5CryptoServiceProvider();
- byte[] data = Encoding.UTF8.GetBytes(str);
- byte[] result = md5.ComputeHash(data);
- String md5Str = BitConverter.ToString(result).Replace("-","").ToLower();
- return md5Str;
- }
-
-
- class Result
- {
- public String SUCCESS { get; set; }
- public String MSG { get; set; }
- }
- }
- }
-
-
-
|