RouteRfidKafkaAction.class.php 26 KB

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