Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Add SQSServiceSensor for non-polling SQS sensor - #91

Open
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor
Open

Add SQSServiceSensor for non-polling SQS sensor#91
grepory wants to merge 1 commit into
StackStorm-Exchange:masterfrom
grepory:grepory/non-polling-sqs-sensor

Conversation

@grepory

Copy link
Copy Markdown

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.

Closes#90 cc @Kami

The SQSServiceSensor class name is meh, but I lack in creativity. This is what we're using right now to process anywhere between 30-100 messages per second from a single queue.

This adds a SQS Sensor with its own polling loop so that we can
consume messages from one or more SQS queues as quickly as possible
without relying on StackStorm to trigger a poll interval.
@grepory
greporyforce-pushed the grepory/non-polling-sqs-sensor branch from 0651675 to 50d0d31CompareDecember 18, 2019 23:51
# setting SQS ServiceResource object from the parameter of datastore or configuration file
self._may_setup_sqs()

while True:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume there is no yielding needed in this function (aka eventlet.sleep(0.01) at the end or similar) because _receive_messages performs a network operation which already needs to yield at some point even if there are no messages to be retrieved.

Otherwise if that's not the case and _receive_messages could immediately return this could cause CPU spikes and 100% CPU utilization by the sensor process since there is no yielding.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

boto3 receive_messages accepts a WaitTimeSeconds argument, which _receive_messages defaults to 2 seconds. That's what's keeping the loop from spinning too fast.

---
class_name: "AWSSQSServiceSensor"
entry_point: "sqs_service_sensor.py"
description: "Service Sensor which monitors a SQS queue for new messages"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would also be good to document here in the description and also in README how this sensor differentiates from other one :)

payload = {"queue": queue, "body": json.loads(msg.body)}
self._sensor_service.dispatch(trigger="aws.sqs_new_message",
payload=payload)
msg.delete()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this method throw?

If so, it probably wouldn't be a bad idea to wrap it in try / catch to avoid a scenario where the same message would always throw for some reason which would prevent sensor from continuing the processing since it would always crash on exception...

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If this throws an exception, it's not because SQS couldn't delete the message due to some unrecoverable error. It will be because of a misconfiguration on the client side or something like that. I don't think it is necessary to try/catch this.

self._logger.warning("SQS Queue: %s doesn't exist, creating it.", queueName)
return self.sqs_res.create_queue(QueueName=queueName)
elif e.response['Error']['Code'] == 'InvalidClientTokenId':
self._logger.warning("Cloudn't operate sqs because of invalid credential config")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also throw (and abort sensor processing) on this error?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know why this was logged and not raised. It shouldn't be possible to get here if you have invalid credentials. We tested that scenario, and it raises an exception. I can't recall if it did so when getting the resource or creating the session.

@punkrokk

Copy link
Copy Markdown
Contributor

@grepory Hey - I want to get this merged, so I am reviewing this PR and I have a few questions.

  1. What happens if someone adds a large amount of queues? Say 30? 100? 500?
  2. Any system requirements recommendations? I can see this sensor getting people into trouble without care consideration before turning it on.
  3. What if you named the sensor SQSContinuousSensor or something like that to better differentiate?

@grepory

grepory commented Feb 8, 2020 via email

Copy link
Copy Markdown
Author

@CLAassistant

CLAassistant commented May 11, 2022

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.


Greg Poirier seems not to be a GitHub user. You need a GitHub account to be able to sign the CLA. If you have already a GitHub account, please add the email address used for this commit to your account.
You have signed the CLA already but the status is still pending? Let us recheck it.

@satellite-no

Copy link
Copy Markdown
Contributor

Based on @grepory last comment shouldnt this be closed?

Can we implement Stale-bot to clean old PRs issues?

@alaypatel07

alaypatel07 commented Feb 14, 2024

Copy link
Copy Markdown

I am working on a usecase that needs the SQS sensor to churn at a faster rate than 120 messages/minute, which is the limitation of current in-tree sensor.

Considering that I would be interested in moving this PR forward. Also, the non-polling sensor could potentially benefit from an async mechanism to grab the sqs message. Would exploring the following https://aiobotocore.readthedocs.io/en/latest/examples/sqs/producer_consumer.html# for non polling sqs sensor help address the open items on the review esp regarding timeout interval?

cc @Kami ^ since you were helping with the review

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Change to Sensor from PollingSensor for SQS Sensor

6 participants

@grepory@punkrokk@CLAassistant@satellite-no@alaypatel07@Kami