Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
S
services
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
tencent
services
Commits
cbc8b72f
Commit
cbc8b72f
authored
Feb 11, 2019
by
xmy
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
feat:配置转移
主题发送日志修改
parent
6068f312
Hide whitespace changes
Inline
Side-by-side
Showing
4 changed files
with
293 additions
and
277 deletions
+293
-277
src/Cmq/Topic.php
+272
-248
src/Common/Config/CfgCenter.php
+1
-1
src/Common/Config/config.php
+16
-26
src/Common/Lib/Xcrypt.php
+4
-2
No files found.
src/Cmq/Topic.php
View file @
cbc8b72f
...
...
@@ -2,245 +2,248 @@
namespace
Hdll\Services\Cmq
;
use
Hdll\Services\Common\Config\CfgCenter
;
use
Swoft\App
;
use
Hdll\Services\Common\Lib\Xcrypt
;
class
Topic
{
private
$topic_name
;
private
$cmq_client
;
private
$encoding
;
public
function
__construct
(
$topic_name
,
$cmq_client
,
$encoding
=
false
)
{
$this
->
topic_name
=
$topic_name
;
$this
->
cmq_client
=
$cmq_client
;
$this
->
encoding
=
$encoding
;
}
public
function
set_encoding
(
$encoding
)
{
$this
->
encoding
=
$encoding
;
}
/*
* create topic
* @type topic_meta : TopicMeta
* @param topic_meta :
*/
public
function
create
(
$topic_meta
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'filterType'
=>
$topic_meta
->
filterType
,
);
if
(
$topic_meta
->
maxMsgSize
>
0
)
{
$params
[
'maxMsgSize'
]
=
$topic_meta
->
maxMsgSize
;
}
$this
->
cmq_client
->
create_topic
(
$params
);
}
/*
* get attributes
*
* @return topic_meta :TopicMeta
*
*/
public
function
get_attributes
()
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
$resp
=
$this
->
cmq_client
->
get_topic_attributes
(
$params
);
$topic_meta
=
new
TopicMeta
();
$this
->
__resp2meta
(
$topic_meta
,
$resp
);
return
$topic_meta
;
}
/*
* set attributes
*
* @type topic_meta :TopicMeta
* @param topic_meta :
*/
public
function
set_attributes
(
$topic_meta
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'maxMsgSize'
=>
strval
(
$topic_meta
->
maxMsgSize
),
);
$this
->
cmq_client
->
set_topic_attributes
(
$params
);
}
/*
* delete topic
*/
public
function
delete
()
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
$this
->
cmq_client
->
delete_topic
(
$params
);
}
/*
* 推送消息 非批量
* @type message :string
* @param message
*
* @type vTagList :list
* @param vTagList 标签
*
* @return message handle
*/
public
function
publish_message
(
$message
,
$vTagList
=
null
,
$routingKey
=
null
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'msgBody'
=>
$message
,
);
if
(
$routingKey
!=
null
)
{
$params
[
'routingKey'
]
=
$routingKey
;
}
if
(
$vTagList
!=
null
&&
is_array
(
$vTagList
)
&&
!
empty
(
$vTagList
))
{
$n
=
1
;
foreach
(
$vTagList
as
$tag
)
{
$key
=
'msgTag.'
.
$n
;
$params
[
$key
]
=
$tag
;
$n
+=
1
;
}
}
$msgId
=
$this
->
cmq_client
->
publish_message
(
$params
);
return
$msgId
;
}
/*
* 批量推送消息
* @type vmessageList :list
* @param vmessageList:
*
* @type vtagList :list
* @param vtagList
*
* @return : return message handle list
*/
public
function
batch_publish_message
(
$vmessageList
,
$vtagList
=
null
,
$routingKey
=
null
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
if
(
$routingKey
!=
null
)
{
$params
[
'routingKey'
]
=
$routingKey
;
}
$n
=
1
;
if
(
is_array
(
$vmessageList
)
&&
!
empty
(
$vmessageList
))
{
foreach
(
$vmessageList
as
$msg
)
{
$key
=
'msgBody.'
.
$n
;
if
(
$this
->
encoding
)
{
$params
[
$key
]
=
base64_encode
(
$msg
);
}
else
{
$params
[
$key
]
=
$msg
;
}
$n
+=
1
;
}
}
if
(
$vtagList
!=
null
&&
is_array
(
$vtagList
)
&&
!
empty
(
$vtagList
))
{
$n
=
1
;
foreach
(
$vtagList
as
$tag
)
{
$key
=
'msgTag.'
.
$n
;
$params
[
$key
]
=
$tag
;
$n
+=
1
;
}
}
$msgList
=
$this
->
cmq_client
->
batch_publish_message
(
$params
);
$retMessageList
=
array
();
foreach
(
$msgList
as
$msg
)
{
if
(
isset
(
$msg
[
'msgId'
]))
{
$retmsgId
=
$msg
[
'msgId'
];
$retMessageList
[]
=
$retmsgId
;
}
}
return
$retMessageList
;
}
/* 列出Topic的Subscriptoin
@type topic_name :string
@param topic_name:
@type searchWord: string
@param searchWord: 订阅关键字
@type limit: int
@param limit: 最多返回的订阅数目
@type offset: string
@param offset: list_subscription的起始位置,上次list_subscription返回的next_offset
@rtype: tuple
@return: subscriptionURL的列表和下次list subscription的起始位置; 如果所有subscription都list出来,next_offset为"".
*/
public
function
list_subscription
(
$searchWord
=
""
,
$limit
=
-
1
,
$offset
=
""
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
);
if
(
$searchWord
!=
""
)
{
$params
[
'searchWord'
]
=
$searchWord
;
}
if
(
$limit
!=
-
1
)
{
$params
[
'limit'
]
=
$limit
;
}
if
(
$offset
!=
""
)
{
$params
[
'offset'
]
=
$offset
;
}
$resp
=
$this
->
cmq_client
->
list_subscription
(
$params
);
if
(
$offset
==
""
)
{
$next_offset
=
count
(
$resp
[
'subscriptionList'
]);
}
else
{
$next_offset
=
$offset
+
count
(
$resp
[
'subscriptionList'
]);
}
if
(
$next_offset
>=
$resp
[
'totalCount'
])
{
$next_offset
=
""
;
}
return
array
(
"totalCoult"
=>
$resp
[
'totalCount'
],
"subscriptionList"
=>
$resp
[
'subscriptionList'
],
"next_offset"
=>
$next_offset
);
}
protected
function
__resp2meta
(
$topic_meta
,
$resp
)
{
if
(
isset
(
$resp
[
'maxMsgSize'
]))
{
$topic_meta
->
maxMsgSize
=
$resp
[
'maxMsgSize'
];
}
if
(
isset
(
$resp
[
'msgRetentionSeconds'
]))
{
$topic_meta
->
msgRetentionSeconds
=
$resp
[
'msgRetentionSeconds'
];
}
if
(
isset
(
$resp
[
'createTime'
]))
{
$topic_meta
->
createTime
=
$resp
[
'createTime'
];
}
if
(
isset
(
$resp
[
'lastModifyTime'
]))
{
$topic_meta
->
lastModifyTime
=
$resp
[
'lastModifyTime'
];
}
if
(
isset
(
$resp
[
'filterType'
]))
{
$topic_meta
->
filterType
=
$resp
[
'filterType'
];
}
}
private
$topic_name
;
private
$cmq_client
;
private
$encoding
;
public
function
__construct
(
$topic_name
,
$cmq_client
,
$encoding
=
false
)
{
$this
->
topic_name
=
$topic_name
;
$this
->
cmq_client
=
$cmq_client
;
$this
->
encoding
=
$encoding
;
}
public
function
set_encoding
(
$encoding
)
{
$this
->
encoding
=
$encoding
;
}
/*
* create topic
* @type topic_meta : TopicMeta
* @param topic_meta :
*/
public
function
create
(
$topic_meta
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'filterType'
=>
$topic_meta
->
filterType
,
);
if
(
$topic_meta
->
maxMsgSize
>
0
)
{
$params
[
'maxMsgSize'
]
=
$topic_meta
->
maxMsgSize
;
}
$this
->
cmq_client
->
create_topic
(
$params
);
}
/*
* get attributes
*
* @return topic_meta :TopicMeta
*
*/
public
function
get_attributes
()
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
$resp
=
$this
->
cmq_client
->
get_topic_attributes
(
$params
);
$topic_meta
=
new
TopicMeta
();
$this
->
__resp2meta
(
$topic_meta
,
$resp
);
return
$topic_meta
;
}
/*
* set attributes
*
* @type topic_meta :TopicMeta
* @param topic_meta :
*/
public
function
set_attributes
(
$topic_meta
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'maxMsgSize'
=>
strval
(
$topic_meta
->
maxMsgSize
),
);
$this
->
cmq_client
->
set_topic_attributes
(
$params
);
}
/*
* delete topic
*/
public
function
delete
()
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
$this
->
cmq_client
->
delete_topic
(
$params
);
}
/*
* 推送消息 非批量
* @type message :string
* @param message
*
* @type vTagList :list
* @param vTagList 标签
*
* @return message handle
*/
public
function
publish_message
(
$message
,
$vTagList
=
null
,
$routingKey
=
null
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
'msgBody'
=>
$message
,
);
if
(
$routingKey
!=
null
)
{
$params
[
'routingKey'
]
=
$routingKey
;
}
if
(
$vTagList
!=
null
&&
is_array
(
$vTagList
)
&&
!
empty
(
$vTagList
))
{
$n
=
1
;
foreach
(
$vTagList
as
$tag
)
{
$key
=
'msgTag.'
.
$n
;
$params
[
$key
]
=
$tag
;
$n
+=
1
;
}
}
$msgId
=
$this
->
cmq_client
->
publish_message
(
$params
);
return
$msgId
;
}
/*
* 批量推送消息
* @type vmessageList :list
* @param vmessageList:
*
* @type vtagList :list
* @param vtagList
*
* @return : return message handle list
*/
public
function
batch_publish_message
(
$vmessageList
,
$vtagList
=
null
,
$routingKey
=
null
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
,
);
if
(
$routingKey
!=
null
)
{
$params
[
'routingKey'
]
=
$routingKey
;
}
$n
=
1
;
if
(
is_array
(
$vmessageList
)
&&
!
empty
(
$vmessageList
))
{
foreach
(
$vmessageList
as
$msg
)
{
$key
=
'msgBody.'
.
$n
;
if
(
$this
->
encoding
)
{
$params
[
$key
]
=
base64_encode
(
$msg
);
}
else
{
$params
[
$key
]
=
$msg
;
}
$n
+=
1
;
}
}
if
(
$vtagList
!=
null
&&
is_array
(
$vtagList
)
&&
!
empty
(
$vtagList
))
{
$n
=
1
;
foreach
(
$vtagList
as
$tag
)
{
$key
=
'msgTag.'
.
$n
;
$params
[
$key
]
=
$tag
;
$n
+=
1
;
}
}
$msgList
=
$this
->
cmq_client
->
batch_publish_message
(
$params
);
$retMessageList
=
array
();
foreach
(
$msgList
as
$msg
)
{
if
(
isset
(
$msg
[
'msgId'
]))
{
$retmsgId
=
$msg
[
'msgId'
];
$retMessageList
[]
=
$retmsgId
;
}
}
return
$retMessageList
;
}
/* 列出Topic的Subscriptoin
@type topic_name :string
@param topic_name:
@type searchWord: string
@param searchWord: 订阅关键字
@type limit: int
@param limit: 最多返回的订阅数目
@type offset: string
@param offset: list_subscription的起始位置,上次list_subscription返回的next_offset
@rtype: tuple
@return: subscriptionURL的列表和下次list subscription的起始位置; 如果所有subscription都list出来,next_offset为"".
*/
public
function
list_subscription
(
$searchWord
=
""
,
$limit
=
-
1
,
$offset
=
""
)
{
$params
=
array
(
'topicName'
=>
$this
->
topic_name
);
if
(
$searchWord
!=
""
)
{
$params
[
'searchWord'
]
=
$searchWord
;
}
if
(
$limit
!=
-
1
)
{
$params
[
'limit'
]
=
$limit
;
}
if
(
$offset
!=
""
)
{
$params
[
'offset'
]
=
$offset
;
}
$resp
=
$this
->
cmq_client
->
list_subscription
(
$params
);
if
(
$offset
==
""
)
{
$next_offset
=
count
(
$resp
[
'subscriptionList'
]);
}
else
{
$next_offset
=
$offset
+
count
(
$resp
[
'subscriptionList'
]);
}
if
(
$next_offset
>=
$resp
[
'totalCount'
])
{
$next_offset
=
""
;
}
return
array
(
"totalCoult"
=>
$resp
[
'totalCount'
],
"subscriptionList"
=>
$resp
[
'subscriptionList'
],
"next_offset"
=>
$next_offset
);
}
protected
function
__resp2meta
(
$topic_meta
,
$resp
)
{
if
(
isset
(
$resp
[
'maxMsgSize'
]))
{
$topic_meta
->
maxMsgSize
=
$resp
[
'maxMsgSize'
];
}
if
(
isset
(
$resp
[
'msgRetentionSeconds'
]))
{
$topic_meta
->
msgRetentionSeconds
=
$resp
[
'msgRetentionSeconds'
];
}
if
(
isset
(
$resp
[
'createTime'
]))
{
$topic_meta
->
createTime
=
$resp
[
'createTime'
];
}
if
(
isset
(
$resp
[
'lastModifyTime'
]))
{
$topic_meta
->
lastModifyTime
=
$resp
[
'lastModifyTime'
];
}
if
(
isset
(
$resp
[
'filterType'
]))
{
$topic_meta
->
filterType
=
$resp
[
'filterType'
];
}
}
/**
* 加密 发送主题消息
...
...
@@ -250,18 +253,39 @@ class Topic
* @return mixed
* @author work
*/
public
function
cryptPushMessage
(
string
$message
,
$vTagList
=
null
,
$routingKey
=
null
){
$cryptMessage
=
Xcrypt
::
encrypt
(
$message
);
$tryTimes
=
0
;
do
{
$res
=
$this
->
publish_message
(
$cryptMessage
,
$vTagList
,
$routingKey
);
$tryTimes
++
;
}
while
(
$res
[
'code'
]
!=
0
&&
$tryTimes
<
3
);
if
(
$tryTimes
>=
3
){
App
::
error
(
"[消息队列失败]:
$message
"
);
}
return
$res
;
}
public
function
cryptPushMessage
(
string
$message
,
$vTagList
=
null
,
$routingKey
=
null
)
{
$cryptMessage
=
Xcrypt
::
encrypt
(
$message
);
$tryTimes
=
0
;
do
{
$res
=
$this
->
publish_message
(
$cryptMessage
,
$vTagList
,
$routingKey
);
$tryTimes
++
;
}
while
(
$res
[
'code'
]
!=
0
&&
$tryTimes
<
3
);
if
(
$tryTimes
>=
3
)
{
App
::
error
(
"[消息队列失败]:
$message
"
);
}
$this
->
TopicLog
(
$message
,
$cryptMessage
,
$vTagList
,
$res
);
return
$res
;
}
protected
function
TopicLog
(
$message
,
$cryptMessage
,
$tagName
,
$response
)
{
$data
=
[
'tagName'
=>
$tagName
,
'topName'
=>
$this
->
topic_name
,
'message'
=>
$message
,
'cryptMessage'
=>
$cryptMessage
,
'createTime'
=>
time
(),
'response'
=>
json_encode
(
$response
),
];
try
{
$db
=
CfgCenter
::
dbConnect
();
$db
->
insert
(
'topic_log'
,
$data
);
}
catch
(
\Exception
$e
){
App
::
error
(
"消息主题日志记录失败:"
.
$e
->
getMessage
()
.
'---'
.
json_encode
(
$data
));
}
}
}
src/Common/Config/CfgCenter.php
View file @
cbc8b72f
...
...
@@ -68,7 +68,7 @@ class CfgCenter
return
[
trim
(
$rkey
,
':'
),
$valObj
];
}
p
rotected
static
function
dbConnect
()
p
ublic
static
function
dbConnect
()
{
if
(
\env
(
'ENVIRONMENT'
,
''
)
==
''
)
{
// 返回线上数据库连接
...
...
src/Common/Config/config.php
View file @
cbc8b72f
<?php
return
[
'qCloud'
=>
[
'Bucket'
=>
'hdll-1257143824'
,
'APPID'
=>
'1257143824'
,
'SecretId'
=>
'AKIDseHj18kua0KTSJ4g9SadbVEnEUZVjvPj'
,
'SecretKey'
=>
'IPL5g5PaaSAzd6NSO8gEmLxcN4pTzJSQ'
,
'Region'
=>
'ap-shanghai'
],
'alisms'
=>
[
'accessKeyId'
=>
'EjBn9zQxyEkKHyAA'
,
'accessKeySecret'
=>
'AN276rwCcqCkFUVt1GLCbAy8jnj52t'
,
],
'cls'
=>
[
'appid'
=>
'1257143824 '
,
'secretId'
=>
'AKIDseHj18kua0KTSJ4g9SadbVEnEUZVjvPj'
,
'secretKey'
=>
'IPL5g5PaaSAzd6NSO8gEmLxcN4pTzJSQ'
,
],
'cmq'
=>
[
'intranet_host'
=>
'http://cmq-topic-bj.api.tencentyun.com'
,
//内网
'internet_host'
=>
'https://cmq-topic-bj.api.qcloud.com'
,
//外网
'secretId'
=>
'AKIDYRW1cG2iVIg8dAoCe86vhNuA7A5oNknk'
,
'secretKey'
=>
'z0ymS7xfLrP6Sk2EKHWaXN6d0EIxX0IQ'
,
],
'cryptKey'
=>
'ebf032f01aa2093be3ee2ee2c137hdll'
,
];
\ No newline at end of file
namespace
Hdll\Services\Common\Config
;
use
Swoft\Redis\Redis
;
use
Swoft\App
;
const
CACHE_PREFIX
=
'CONFIG_CACHE'
;
$redis
=
App
::
getBen
(
Redis
::
class
);
$data
=
$redis
->
get
(
CACHE_PREFIX
.
'all'
);
if
(
!
empty
(
$data
))
{
$db
=
CfgCenter
::
dbConnect
();
$data
=
$db
->
select
(
'config'
,
[
'name'
,
'value'
]);
$redis
->
set
(
CACHE_PREFIX
.
'all'
,
$data
);
}
return
$data
;
\ No newline at end of file
src/Common/Lib/Xcrypt.php
View file @
cbc8b72f
<?php
namespace
Hdll\Services\Common\Lib
;
use
Hdll\Services\Common\Config\CfgCenter
;
use
Swoft\App
;
use
Swoft\Redis\Redis
;
...
...
@@ -41,8 +42,9 @@ class Xcrypt
}
private
static
function
getKey
(){
$redis
=
App
::
getBean
(
Redis
::
class
);
$key
=
$redis
->
get
(
self
::
CRYPT
);
// $redis = App::getBean(Redis::class);
// $key = $redis->get(self::CRYPT);
$key
=
CfgCenter
::
get
(
self
::
CRYPT
);
if
(
empty
(
$key
)){
throw
new
\Exception
(
'加密密钥获取失败!'
);
}
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment