RouteRfidKafkaAction.class.php 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876
  1. <?php
  2. class RouteRfidKafkaAction extends Action {
  3. private function addRfidDataToNingbo( $data, $conn ){
  4. //{"methond":"track","mac":"ffc10063","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"01000423","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621112931"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621113202"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"160621113303"}]}
  5. $this->addMonitorProcess();
  6. if($data['methond']=='login'){
  7. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond login!');
  8. }
  9. $RF_ID=strtoupper($data['mac']);
  10. $station_cond=array('mac'=>$RF_ID);
  11. $device_name=M('stations')->where($station_cond)->getField('name');
  12. if(!$device_name){
  13. //return array('success'=>true,'message'=>'addRfidDataToNingbo failed,station not existed!');
  14. }
  15. if($data['methond']=='heartbeat'){
  16. if( ($data['time']<(time()-3600) ) || ($data['time']>(time()+3600) ) ){
  17. $this->debug_log( 'heartbeat_abnormal', $data );
  18. return array('success'=>true,'message'=>'heartbeat time abnormal !');
  19. }
  20. $save_data=array(
  21. 'online_time'=>date('Y-m-d H:i:s',$data['time'])
  22. );
  23. M('stations')->createSave($station_cond,$save_data);
  24. }
  25. if($data['methond']!='track'){
  26. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond not track!');
  27. }
  28. if(!$data['labels']){
  29. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,labels not existed!');
  30. }
  31. $save_data=array(
  32. 'online_time'=>date('Y-m-d H:i:s',time())
  33. );
  34. M('stations')->createSave($station_cond,$save_data);
  35. //$station_sql='SELECT DEVICE_NAME FROM DSSC2.ADM_DEV WHERE LOGIN_NAME=\''.$RF_ID.'\'';
  36. //$stid = oci_parse($conn, $station_sql);
  37. //oci_define_by_name($stid, 'DEVICE_NAME', $device_name);
  38. //oci_execute($stid);
  39. //oci_fetch($stid);
  40. foreach($data['labels'] as $val){
  41. if( ($val['time']<(time()-3600*24*10) ) || ($val['time']>(time()+3600*24*10) ) ){
  42. $val['mac']=$RF_ID;
  43. $this->debug_log( 'abnormal_labels', $val );
  44. continue;
  45. }
  46. $RF_STAT=0;
  47. $plate_no='';
  48. if($val['event']['entry']==1){
  49. $RF_STAT=1;
  50. }elseif($val['event']['leave']==1){
  51. $RF_STAT=2;
  52. }
  53. $RF_FLAGID=strtoupper($val['id']);
  54. if($RF_FLAGID=='00000000'){
  55. $val['mac']=$RF_ID;
  56. $this->debug_log( 'abnormal_labels', $val);
  57. continue;
  58. }
  59. if(!$val['time']){
  60. continue;
  61. }
  62. $RF_DATE=date('Y-m-d H:i:s',$val['time']);
  63. $sql = 'INSERT INTO "DSSC2"."W_DW_RF_RECORD"("ID", "RF_ID", "RF_FLAGID", "RF_DATE", "RF_STAT") VALUES (DSSC2.SEQ_W_DW_RF_RECORD.nextval, \''.$RF_ID.'\', \''.$RF_FLAGID.'\', TO_DATE(\''.$RF_DATE.'\', \'SYYYY-MM-DD HH24:MI:SS\'), \''.$RF_STAT.'\')';
  64. //var_dump($sql);
  65. //$station_str=$RF_ID.' '.$RF_FLAGID.' '.$RF_DATE;
  66. //$this->debug_log( 'station_lebal_data', $station_str);
  67. //插入数据到oracle轨迹表
  68. $stid = oci_parse($conn, $sql);
  69. $r = oci_execute($stid);
  70. if(!$r){
  71. $val['mac']=$RF_ID;
  72. $this->debug_log( 'insert_oracle_error', $val );
  73. return array('success'=>false,'message'=>'addRfidDataToNingbo failed,insert_oracle_error!');
  74. }
  75. if(!$device_name){
  76. continue;
  77. }
  78. //插入成功就执行统计
  79. $handle_data=array(
  80. 'RF_STAT'=>$RF_STAT,
  81. 'RF_FLAGID'=>$RF_FLAGID,
  82. 'RF_ID'=>$RF_ID,
  83. 'time'=>$val['time'],
  84. 'address'=>$device_name,
  85. );
  86. $this->handleTotalData($handle_data);
  87. $vehicle_sql='SELECT o.PLATE_NO FROM DSSC3.W_DW_NON_MOTOR o,DSSC3.W_DW_RFID_TAGS s WHERE s.RFID_SN =\''.$RF_FLAGID.'\' AND o.rfid_id = s.id ';
  88. $stid = oci_parse($conn, $vehicle_sql);
  89. oci_define_by_name($stid, 'PLATE_NO', $plate_no);
  90. oci_execute($stid);
  91. oci_fetch($stid);
  92. if(!$plate_no){
  93. continue;
  94. }
  95. $handle_data['plate_no']=$plate_no;
  96. //检测布控
  97. $this->checkControlAlarm($handle_data);
  98. //违规行驶检测 超速逆行检测
  99. $this->checkIllegalDriving($handle_data);
  100. }
  101. oci_free_statement($stid);
  102. //oci_close($conn);
  103. return array('success'=>true,'message'=>'add success');
  104. }
  105. public function pushRfidRouteToNingbo( ){
  106. $broker_list = C('KAFKA_BROKER_LIST');
  107. if (empty($broker_list)) {
  108. exit("KAFKA_BROKER_LIST must be config!".PHP_EOL);
  109. }
  110. $group = C('ROUTE_INDEX_KAFKA_GROUP');
  111. if (empty($group)) {
  112. exit("ROUTE_INDEX_KAFKA_GROUP must be config!".PHP_EOL);
  113. }
  114. $topics = C('ROUTE_INDEX_KAFKA_TOPIC');
  115. if (empty($topics)) {
  116. exit("ROUTE_INDEX_KAFKA_TOPIC must be config!".PHP_EOL);
  117. }
  118. $topics = explode(',',$topics);
  119. // 从 topic :rlstation_rfid_location 取轨迹
  120. $conf = new RdKafka\Conf();
  121. // Set a rebalance callback to log partition assignments (optional)(当有新的消费进程加入或者退出消费组时,kafka 会自动重新分配分区给消费者进程,这里注册了一个回调函数,当分区被重新分配时触发)
  122. $conf->setRebalanceCb(function (RdKafka\KafkaConsumer $kafka, $err, array $partitions = null) {
  123. switch ($err) {
  124. case RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS:
  125. echo "Assign: ";
  126. var_dump($partitions);
  127. $kafka->assign($partitions);
  128. break;
  129. case RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS:
  130. echo "Revoke: ";
  131. var_dump($partitions);
  132. $kafka->assign(NULL);
  133. break;
  134. default:
  135. throw new \Exception($err);
  136. }
  137. });
  138. // Configure the group.id. All consumer with the same group.id will consume( 配置groud.id 具有相同 group.id 的consumer 将会处理不同分区的消息,所以同一个组内的消费者数量如果订阅了一个topic, 那么消费者进程的数量多于这个topic 分区的数量是没有意义的。)
  139. // different partitions.
  140. $conf->set('group.id', $group);
  141. // Initial list of Kafka brokers(添加 kafka集群服务器地址)
  142. $conf->set('metadata.broker.list', $broker_list);
  143. $topicConf = new RdKafka\TopicConf();
  144. // Set where to start consuming messages when there is no initial offset in
  145. // offset store or the desired offset is out of range.
  146. // 'smallest': start from the beginning
  147. $topicConf->set('auto.offset.reset', 'latest');
  148. // Set the configuration to use for subscribed/assigned topics
  149. $conf->setDefaultTopicConf($topicConf);
  150. $consumer = new RdKafka\KafkaConsumer($conf);
  151. // 订阅轨迹数据topic
  152. $consumer->subscribe($topics);
  153. $config = C('ORACLE_CONFIG');
  154. if (empty($config)) {
  155. exit("ORACLE_CONFIG must be config!".PHP_EOL);
  156. }
  157. $host= $config['host'];
  158. $port= $config['port'];
  159. $instance_name= $config['instance_name'];
  160. $username= $config['username'];
  161. $password= $config['password'];
  162. /*
  163. $host= '192.168.100.23';
  164. $port= '1521';
  165. $instance_name= 'helowin';
  166. $username= 'DSSC3';
  167. $password= 'Rliandssc3';
  168. */
  169. $conn = oci_connect($username, $password, $host.':'.$port.'/'. $instance_name,'AL32UTF8');
  170. if (!$conn) {
  171. $e = oci_error();
  172. trigger_error(htmlentities($e['message'], ENT_QUOTES), E_USER_ERROR);
  173. }
  174. while (true) {
  175. //var_dump($conn);
  176. $message = $consumer->consume(120*1000);
  177. switch ($message->err) {
  178. case RD_KAFKA_RESP_ERR_NO_ERROR:
  179. $data = json_decode($message->payload,true);
  180. if( $data ){
  181. $res=$this->addRfidDataToNingbo($data,$conn);
  182. if(!$res['success']){
  183. throw new \Exception($res['message']);
  184. }
  185. //$this->addRfidDataToRenlian($data);
  186. }
  187. break;
  188. case RD_KAFKA_RESP_ERR__PARTITION_EOF:
  189. //echo "No more messages; will wait for more\n";
  190. break;
  191. case RD_KAFKA_RESP_ERR__TIMED_OUT:
  192. echo "Timed out\n";
  193. break;
  194. default:
  195. echo "default break";
  196. $this->debug_log( 'default_Log', $message->errstr() );
  197. $this->debug_log( 'default_Log', $message->err );
  198. break;
  199. }
  200. }
  201. }
  202. private function addRfidDataToRenlian( $data ){
  203. //{"methond":"track","mac":"ffc10063","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"01000423","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621112931"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621113202"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"160621113303"}]}
  204. //var_dump($data);
  205. if($data['methond']!='track'){
  206. return array('success'=>false,'message'=>'addRfidDataToNingbo failed,methond not track!');
  207. }
  208. if(!$data['labels']){
  209. return array('success'=>false,'message'=>'addRfidDataToNingbo failed,labels not existed!');
  210. }
  211. var_dump($data);
  212. $conn = null;
  213. $host= '115.198.203.63';
  214. $port= '1521';
  215. $instance_name= 'ORCL';
  216. $username= 'DSSC3';
  217. $password= '123456';
  218. $conn = new PDO("oci:dbname=//".$host.":".$port."/".$instance_name,$username,$password);// PDO方式
  219. $conn->setAttribute(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION);
  220. $RF_ID=strtoupper($data['mac']);
  221. foreach($data['labels'] as $val){
  222. $RF_STAT=0;
  223. if($val['event']['entry']==1){
  224. $RF_STAT=1;
  225. }elseif($val['event']['leave']==1){
  226. $RF_STAT=2;
  227. }
  228. $RF_FLAGID=strtoupper($val['id']);
  229. $RF_DATE=date('Y-m-d H:i:s',$val['time']);
  230. $sql = 'INSERT INTO "ROOT"."W_DW_RF_RECORD"("RF_ID", "RF_FLAGID", "RF_DATE", "RF_STAT") VALUES (\''.$RF_ID.'\', \''.$RF_FLAGID.'\', TO_DATE(\''.$RF_DATE.'\', \'SYYYY-MM-DD HH24:MI:SS\'), \''.$RF_STAT.'\')';
  231. var_dump($recvCount . '.' . $sql);
  232. // $res = $conn -> query($sql);
  233. }
  234. if ($conn){
  235. $conn = null;
  236. }
  237. return array('success'=>true,'message'=>'add success');
  238. }
  239. private function addControlAlarm( $data, $params ){
  240. $save_data=array(
  241. 'plate_no'=>$data['plate_no'],
  242. 'rfid_sn'=>$data['RF_FLAGID'],
  243. 'address'=>$data['address'],
  244. 'alarm_type'=>$data['alarm_type'],
  245. 'created_at'=>$data['time'],
  246. 'remark'=>$data['remark'],
  247. 'state'=>1,
  248. );
  249. $res=M('control_alarm')->createAdd($save_data);
  250. if($res && $params['is_sent'] && $params['phone']){
  251. $time=time()-$params['sent_limit'];
  252. $where=[
  253. 'recipient'=>$params['phone'],
  254. 'type'=>$data['alarm_type'],
  255. 'created_at'=>array('GT',$time)
  256. ];
  257. $count=M('sms_send_log')->where($where)->count();
  258. if($count){
  259. return $res;
  260. }
  261. $sent_data = [
  262. "signName" => "任联科技",
  263. "templateCode" => "SMS_273620690",
  264. "phoneNumbers" => $params['phone'],
  265. "templateParam" => json_encode(['name' => $data['plate_no'],'time'=>date('Y-m-d H:i:s',$data['time']),'content'=>'在'.$data['address'].'处'.$data['remark']])
  266. ];
  267. $response=sendSmsAli($sent_data);
  268. $sms_save_data=[
  269. 'recipient'=>$params['phone'],
  270. 'type'=>$data['alarm_type'],
  271. 'created_at'=>time(),
  272. ];
  273. if($response['code']=='ok'){
  274. $sms_save_data['sent_result']='发送成功';
  275. }else{
  276. $sms_save_data['sent_result']=$response['message'];
  277. }
  278. M('sms_send_log')->createAdd($sms_save_data);
  279. }
  280. return $res;
  281. }
  282. private function checkControlAlarm( $data ){
  283. //先检测标签是否布控
  284. $cond=array(
  285. 'control_obj'=>array('in',[$data['RF_FLAGID'],$data['plate_no']]),
  286. 'state'=>1
  287. );
  288. $ve_con=M('control_manage')->where($cond)->find();
  289. //存在布控标签 并在时间内
  290. if($ve_con && ($data['time']>$ve_con['start_time']) && ($data['time']<$ve_con['end_time'])){
  291. $data['alarm_type']='control';
  292. $this->addControlAlarm($data,$ve_con);
  293. }
  294. //检测区域布控
  295. $cond2=array(
  296. 'control_obj'=>array("LIKE", '%'.$data['RF_ID'].'%'),
  297. 'state'=>1
  298. );
  299. $sta_con=M('control_manage')->where($cond2)->find();
  300. //存在布控基站
  301. if($sta_con && ($data['time']>$sta_con['start_time']) && ($data['time']<$sta_con['end_time'])){
  302. if($sta_con['bw_ids']){//存在名单
  303. $bwIdArr=explode(',',$sta_con['bw_ids']);
  304. //获取名单内标签
  305. $rfid_arr=M('rfid_with_bw')->where(array('bw_id'=>array('in',$bwIdArr) ))->getField('rfid',true);
  306. if($sta_con['area_type']=='1'){//禁止活动区域
  307. if($sta_con['bw_type']=='0'){//指定黑名单禁止
  308. if(in_array($data['RF_FLAGID'],$rfid_arr)){
  309. //在黑名单内 告警
  310. $data['alarm_type']='forbid_in';
  311. $data['remark']='驶入黑名单禁入区域';
  312. $this->addControlAlarm($data,$sta_con);
  313. }
  314. }else{
  315. //白名单的不禁止
  316. if(!in_array($data['RF_FLAGID'],$rfid_arr)){
  317. //不在白名单内 告警
  318. $data['alarm_type']='forbid_in';
  319. $data['remark']='驶入禁入区域';
  320. $this->addControlAlarm($data,$sta_con);
  321. }
  322. }
  323. }else{//活动区域
  324. //指定黑名单
  325. if(in_array($data['RF_FLAGID'],$rfid_arr)){
  326. //在黑名单内 告警
  327. $data['remark']='驶入活动区域';
  328. $data['alarm_type']='activity_in';
  329. $this->addControlAlarm($data,$sta_con);
  330. }
  331. }
  332. }else{
  333. //无黑白名单布控 全部禁止
  334. $data['remark']='驶入禁入区域';
  335. $data['alarm_type']='forbid_in';
  336. $this->addControlAlarm($data,$sta_con);
  337. }
  338. }
  339. }
  340. private function checkIllegalDriving( $data ){
  341. $redis = Redis("nbfd_stuck_section_data","hash");
  342. //先查询基站是否设置卡点
  343. $cond=array('macs'=>array("LIKE", '%'.$data['RF_ID'].'%'));
  344. $section_id=M('stuck_point')->where($cond)->getField('id');
  345. //检测是否是前置卡点
  346. $pre_section_cond=array('pre_spot'=>$section_id);
  347. $pre_section=M('stuck_section')->where($pre_section_cond)->find();
  348. if($pre_section){
  349. //如果是前置卡点 则记录标签进入卡点区间时间
  350. $key = "stuck_section_".$pre_section['id']."_".$data['RF_FLAGID'];
  351. $passInfo=json_decode($redis->get($key),true);
  352. //存在逆行进入卡点时间 且启用超速检测
  353. if($passInfo && ($passInfo['section']=='pos') && ($pre_section['retrograde_stat']=='1')){
  354. $data['alarm_type']='retrograde';
  355. $data['remark']=$pre_section['name'].'逆行';
  356. $this->addControlAlarm($data);
  357. }
  358. $redisData = array(
  359. $key =>json_encode(array(
  360. "section" => 'pre',
  361. "time" => $data['time'],
  362. )
  363. )
  364. );
  365. $redis->add($redisData);
  366. return;
  367. }
  368. //检查是否是后置卡点
  369. $pos_section_cond=array('pos_spot'=>$section_id);
  370. $pos_section=M('stuck_section')->where($pos_section_cond)->find();
  371. if($pos_section){
  372. //后置卡点 取标签进入卡点区间时间
  373. $key = "stuck_section_".$pos_section['id']."_".$data['RF_FLAGID'];
  374. $passInfo=json_decode($redis->get($key),true);
  375. //存在进入卡点时间 且启用超速检测
  376. if($passInfo && ($passInfo['section']=='pre') && ($pos_section['over_speed_stat']=='1')){
  377. $hour= ($data['time']-$passInfo['time'])/3600;
  378. $speed=($pos_section['distance']/1000)/$hour;
  379. if($speed>$pos_section['max_speed']){
  380. //超速行驶
  381. $data['alarm_type']='over_speed';
  382. $data['remark']=$pos_section['name'].'超速,速度:'. round($speed,2);
  383. $this->addControlAlarm($data);
  384. }
  385. }
  386. //存经过后置卡点时间
  387. $redisData = array(
  388. $key =>json_encode(array(
  389. "section" => 'pos',
  390. "time" => $data['time'],
  391. )
  392. )
  393. );
  394. $redis->add($redisData);
  395. return;
  396. }
  397. }
  398. private function handleTotalData( $data ){
  399. /*
  400. $data=array(
  401. 'RF_STAT'=>$RF_STAT,
  402. 'RF_FLAGID'=>$RF_FLAGID,
  403. 'RF_ID'=>$RF_ID,
  404. 'time'=>$val['time'],
  405. 'address'=>$station_info['DEVICE_NAME']
  406. );
  407. */
  408. //统计表数据添加
  409. $to_cond=array(
  410. 'mac'=>$data['RF_ID'],
  411. 'date'=>date('Y-m-d',$data['time'])
  412. );
  413. if(!M('station_passing')->where($to_cond)->count()){
  414. $total_data=array(
  415. 'address'=>$data['address'],
  416. 'mac'=>$data['RF_ID'],
  417. 'date'=>date('Y-m-d',$data['time']),
  418. 'num'=>1,
  419. );
  420. $res=M('station_passing')->createAdd($total_data);
  421. if(!$res){
  422. $this->debug_log( 'station_passing_error', M('station_passing')->getLastSql());
  423. }
  424. }else{
  425. $res=M('station_passing')->where($to_cond)->setInc('num');
  426. if(!$res){
  427. $this->debug_log( 'station_passing_error', M('station_passing')->getLastSql());
  428. }
  429. }
  430. return $res;
  431. }
  432. public function test( ){
  433. $topic = C('ROUTE_INDEX_KAFKA_TOPIC');
  434. static $rk;
  435. if (!extension_loaded('rdkafka')){
  436. echo 'pushToKafka fail,extension of rdkafka has not installed!!'.PHP_EOL;
  437. return false;
  438. }
  439. if(!$rk){
  440. $conf = new Rdkafka\Conf();
  441. $conf->set('batch.num.messages', 2);
  442. //$conf->set('linger.ms', 10);
  443. //$conf->set('log_level', (string) LOG_DEBUG);
  444. //$conf->set('debug', 'all');
  445. $conf->setErrorCb(function($producer, $msg) {
  446. printf("%s: %s\n", rd_kafka_err2str($err), $errstr);
  447. });
  448. $conf->setDrMsgCb(function($producer, $msg) {
  449. if($msg->err) {
  450. echo 'Message delivery failed:' . $msg->errstr();
  451. } else {
  452. echo "sent message sucessfully.";
  453. }
  454. });
  455. $rk = new RdKafka\Producer($conf);
  456. }
  457. var_dump($topic);
  458. //var_dump(C('KAFKA_BROKER_LIST'));die;
  459. //$rk->setLogLevel(LOG_DEBUG);
  460. $rk->addBrokers(C('KAFKA_BROKER_LIST'));
  461. $topic = $rk->newTopic($topic);
  462. $res='{"methond":"track","mac":"FF0435EE","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"0308EC58","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"'.time().'"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"1670294864"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"1670294864"}]}';
  463. $topic->produce(RD_KAFKA_PARTITION_UA, 0,$res);
  464. $rk->poll(0);
  465. while ($rk->getOutQLen() > 0) {
  466. $rk->poll(1);
  467. }
  468. }
  469. public function test2( ){
  470. $topic = C('ROUTE_INDEX_KAFKA_TOPIC');
  471. static $rk;
  472. if (!extension_loaded('rdkafka')){
  473. echo 'pushToKafka fail,extension of rdkafka has not installed!!'.PHP_EOL;
  474. return false;
  475. }
  476. if(!$rk){
  477. $conf = new Rdkafka\Conf();
  478. $conf->set('batch.num.messages', 2);
  479. //$conf->set('linger.ms', 10);
  480. //$conf->set('log_level', (string) LOG_DEBUG);
  481. //$conf->set('debug', 'all');
  482. $conf->setErrorCb(function($producer, $msg) {
  483. printf("%s: %s\n", rd_kafka_err2str($err), $errstr);
  484. });
  485. $conf->setDrMsgCb(function($producer, $msg) {
  486. if($msg->err) {
  487. echo 'Message delivery failed:' . $msg->errstr();
  488. } else {
  489. echo "sent message sucessfully.";
  490. }
  491. });
  492. $rk = new RdKafka\Producer($conf);
  493. }
  494. var_dump($topic);
  495. //var_dump(C('KAFKA_BROKER_LIST'));die;
  496. //$rk->setLogLevel(LOG_DEBUG);
  497. $rk->addBrokers(C('KAFKA_BROKER_LIST'));
  498. $topic = $rk->newTopic($topic);
  499. $res='{"methond":"track","mac":"FF04C526","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"0308EC58","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"'.time().'"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"1667529922"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"1667529922"}]}';
  500. $info = array(
  501. 'DeviceId' => '869688888888888',
  502. //'State' => (string)$data['state'],
  503. //'Speed' => $data['speed'],
  504. 'Latitude'=>30.192289977186,
  505. 'Longitude'=>120.20063757299,
  506. 'DeviceTime' => time(),
  507. 'Altitude'=>108.951,
  508. //'LBS' => $data['lbs'],
  509. //'Direction' => $data['direction'],
  510. );
  511. $topic->produce(RD_KAFKA_PARTITION_UA, 0,$res);
  512. $rk->poll(0);
  513. while ($rk->getOutQLen() > 0) {
  514. $rk->poll(1);
  515. }
  516. }
  517. public function debug_log( $filename, $data ){
  518. $file = SOLUTION_LOG_PATH .'/'.date("Ymd", time()) ."/".$filename.".log";
  519. $folder=dirname($file);
  520. if (!is_dir($folder)){
  521. mkdir($folder,0777,true);
  522. }
  523. if(is_array($data)){
  524. $data = json_encode($data);
  525. }
  526. file_put_contents($file, '[' . date('Y-m-d H:i:s') . ']' . $data . PHP_EOL,FILE_APPEND);
  527. }
  528. public function test3( ){
  529. $params=[
  530. 'is_sent'=>1,
  531. 'phone'=>'15706857065',
  532. 'sent_limit'=>'600',
  533. ];
  534. if( $params['is_sent'] && $params['phone']){
  535. $time=time()-$params['sent_limit'];
  536. $where=[
  537. 'recipient'=>$params['phone'],
  538. 'type'=>'test',
  539. 'created_at'=>array('GT',$time)
  540. ];
  541. $sent_data = [
  542. "signName" => "任联科技",
  543. "templateCode" => "SMS_273620690",
  544. "phoneNumbers" => $params['phone'],
  545. "templateParam" => json_encode(['name' => '测试A3123','time'=>date('Y-m-d H:i:s',time()),'content'=>'在滨江区华城街道处超速行驶'])
  546. ];
  547. $response=sendSmsAli($sent_data);
  548. $sms_save_data=[
  549. 'recipient'=>$params['phone'],
  550. 'type'=>'test',
  551. 'created_at'=>time(),
  552. ];
  553. if($response['code']=='ok'){
  554. $sms_save_data['sent_result']='发送成功';
  555. }else{
  556. $sms_save_data['sent_result']=$response['message'];
  557. }
  558. var_dump($sms_save_data);
  559. //M('sms_send_log')->createAdd($sms_save_data);
  560. }
  561. return $res;
  562. }
  563. private function addMonitorProcess( ){
  564. $redis = Redis("nbfd_monitor_process_id","hash");
  565. //$pid=getmypid();
  566. $pid=posix_getpid();
  567. $key = "monitor_process_id_".$pid;
  568. $redisData = array(
  569. $key =>json_encode(array(
  570. "pid" => $pid,
  571. "time" => time(),
  572. )
  573. )
  574. );
  575. $redis->add($redisData);
  576. }
  577. public function pushRfidRouteToStudent( ){
  578. $broker_list = C('KAFKA_BROKER_LIST');
  579. if (empty($broker_list)) {
  580. exit("KAFKA_BROKER_LIST must be config!".PHP_EOL);
  581. }
  582. $group = C('ROUTE_INDEX_KAFKA_GROUP_STUDENT');
  583. if (empty($group)) {
  584. exit("ROUTE_INDEX_KAFKA_GROUP_STUDENT must be config!".PHP_EOL);
  585. }
  586. $topics = C('ROUTE_INDEX_KAFKA_TOPIC');
  587. if (empty($topics)) {
  588. exit("ROUTE_INDEX_KAFKA_TOPIC must be config!".PHP_EOL);
  589. }
  590. $topics = explode(',',$topics);
  591. // 从 topic :rlstation_rfid_location 取轨迹
  592. $conf = new RdKafka\Conf();
  593. // Set a rebalance callback to log partition assignments (optional)(当有新的消费进程加入或者退出消费组时,kafka 会自动重新分配分区给消费者进程,这里注册了一个回调函数,当分区被重新分配时触发)
  594. $conf->setRebalanceCb(function (RdKafka\KafkaConsumer $kafka, $err, array $partitions = null) {
  595. switch ($err) {
  596. case RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS:
  597. echo "Assign: ";
  598. var_dump($partitions);
  599. $kafka->assign($partitions);
  600. break;
  601. case RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS:
  602. echo "Revoke: ";
  603. var_dump($partitions);
  604. $kafka->assign(NULL);
  605. break;
  606. default:
  607. throw new \Exception($err);
  608. }
  609. });
  610. // Configure the group.id. All consumer with the same group.id will consume( 配置groud.id 具有相同 group.id 的consumer 将会处理不同分区的消息,所以同一个组内的消费者数量如果订阅了一个topic, 那么消费者进程的数量多于这个topic 分区的数量是没有意义的。)
  611. // different partitions.
  612. $conf->set('group.id', $group);
  613. // Initial list of Kafka brokers(添加 kafka集群服务器地址)
  614. $conf->set('metadata.broker.list', $broker_list);
  615. $topicConf = new RdKafka\TopicConf();
  616. // Set where to start consuming messages when there is no initial offset in
  617. // offset store or the desired offset is out of range.
  618. // 'smallest': start from the beginning
  619. $topicConf->set('auto.offset.reset', 'latest');
  620. // Set the configuration to use for subscribed/assigned topics
  621. $conf->setDefaultTopicConf($topicConf);
  622. $consumer = new RdKafka\KafkaConsumer($conf);
  623. // 订阅轨迹数据topic
  624. $consumer->subscribe($topics);
  625. $config = C('STUDENT_ORACLE_CONFIG');
  626. if (empty($config)) {
  627. exit("STUDENT_ORACLE_CONFIG must be config!".PHP_EOL);
  628. }
  629. $host= $config['host'];
  630. $port= $config['port'];
  631. $instance_name= $config['instance_name'];
  632. $username= $config['username'];
  633. $password= $config['password'];
  634. /*
  635. $host= '192.168.100.23';
  636. $port= '1521';
  637. $instance_name= 'helowin';
  638. $username= 'DSSC3';
  639. $password= 'Rliandssc3';
  640. */
  641. $conn = oci_connect($username, $password, $host.':'.$port.'/'. $instance_name,'AL32UTF8');
  642. if (!$conn) {
  643. $e = oci_error();
  644. trigger_error(htmlentities($e['message'], ENT_QUOTES), E_USER_ERROR);
  645. }
  646. while (true) {
  647. //var_dump($conn);
  648. $message = $consumer->consume(120*1000);
  649. switch ($message->err) {
  650. case RD_KAFKA_RESP_ERR_NO_ERROR:
  651. $data = json_decode($message->payload,true);
  652. if( $data ){
  653. $res=$this->addRfidDataToStudent($data,$conn);
  654. if(!$res['success']){
  655. throw new \Exception($res['message']);
  656. }
  657. }
  658. break;
  659. case RD_KAFKA_RESP_ERR__PARTITION_EOF:
  660. //echo "No more messages; will wait for more\n";
  661. break;
  662. case RD_KAFKA_RESP_ERR__TIMED_OUT:
  663. echo "Timed out\n";
  664. break;
  665. default:
  666. echo "default break";
  667. $this->debug_log( 'default_Log', $message->errstr() );
  668. $this->debug_log( 'default_Log', $message->err );
  669. break;
  670. }
  671. }
  672. }
  673. private function addRfidDataToStudent( $data, $conn ){
  674. //{"methond":"track","mac":"ffc10063","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"01000423","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621112931"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621113202"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"160621113303"}]}
  675. if($data['methond']=='login'){
  676. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond login!');
  677. }
  678. $RF_ID=strtoupper($data['mac']);
  679. $station_cond=array('mac'=>$RF_ID);
  680. $device_name=M('stations')->where($station_cond)->getField('name');
  681. if(!$device_name){
  682. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,station not existed!');
  683. }
  684. if($data['methond']=='heartbeat'){
  685. }
  686. if($data['methond']!='track'){
  687. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond not track!');
  688. }
  689. if(!$data['labels']){
  690. return array('success'=>true,'message'=>'addRfidDataToNingbo failed,labels not existed!');
  691. }
  692. $updateStationTime=true;
  693. foreach($data['labels'] as $val){
  694. if( ($val['time']<(time()-3600*24*10) ) || ($val['time']>(time()+3600*24*10) ) ){
  695. $updateStationTime=false;
  696. $val['mac']=$RF_ID;
  697. continue;
  698. }
  699. $RF_STAT=0;
  700. $plate_no='';
  701. if($val['event']['entry']==1){
  702. $RF_STAT=1;
  703. }elseif($val['event']['leave']==1){
  704. $RF_STAT=2;
  705. }
  706. $RF_FLAGID=strtoupper($val['id']);
  707. if($RF_FLAGID=='00000000'){
  708. continue;
  709. }
  710. $RF_DATE=date('Y-m-d H:i:s',$val['time']);
  711. $sql = 'INSERT INTO "STUDENTCRAD"."W_DW_RF_RECORD"("ID", "RF_ID", "RF_FLAGID", "RF_DATE", "RF_STAT") VALUES (STUDENTCRAD.SEQ_W_DW_RF_RECORD.nextval, \''.$RF_ID.'\', \''.$RF_FLAGID.'\', TO_DATE(\''.$RF_DATE.'\', \'SYYYY-MM-DD HH24:MI:SS\'), \''.$RF_STAT.'\')';
  712. //var_dump($sql);
  713. //插入数据到oracle轨迹表
  714. $stid = oci_parse($conn, $sql);
  715. $r = oci_execute($stid);
  716. if(!$r){
  717. $val['mac']=$RF_ID;
  718. $this->debug_log( 'insert_student_error', $val );
  719. return array('success'=>false,'message'=>'addRfidDataToNingbo failed,insert_oracle_error!');
  720. }
  721. }
  722. oci_free_statement($stid);
  723. return array('success'=>true,'message'=>'add success');
  724. }
  725. }