This message was deleted.
# general
s
This message was deleted.
d
Since this is a Pulsar IO connector, you will have to submit it using the pulsar-admin API for connectors documented here, and in greater detail here.
Copy code
pulsar-admin sinks create --name my-http-sink --archive path/to/download.nar --sink-config-file path/to/your-config.yaml
o
@David K I have tried but I get this error when I check status
Copy code
{
  "numInstances": 1,
  "numRunning": 0,
  "instances": [
    {
      "instanceId": 0,
      "status": {
        "running": false,
        "error": "UNAVAILABLE: io exception",
        "numRestarts": 0,
        "numReadFromPulsar": 0,
        "numSystemExceptions": 0,
        "latestSystemExceptions": [],
        "numSinkExceptions": 0,
        "latestSinkExceptions": [],
        "numWrittenToSink": 0,
        "lastReceivedTime": 0,
        "workerId": "c-standalone-fw-localhost-8080"
      }
    }
  ]
}
d
Can you check the function worker logs for errors? The “UNAVAILABLE: io exception” in the error field looks suspicious. Perhaps the nar file is available on the function-worker node?
o
@David K please do you know the path to the log? because I have been looking for it
d
Is the function-worker running as a K8s pod or a standalone process on a physical server? If the former, then
kubectl logs <POD> --tail
should work. If the latter, I am not sure of the topic of my head.
o
running on standalone
d
Inside docker or from the command line?
o
command line
d
Then the log messages should be in wherever you redirect stdout and stderr to.
o
ok, it should be in logs/functions
@David K thank you very much, I have been able to get it to this point
Copy code
{
  "numInstances": 1,
  "numRunning": 1,
  "instances": [
    {
      "instanceId": 0,
      "status": {
        "running": true,
        "error": "",
        "numRestarts": 0,
        "numReadFromPulsar": 0,
        "numSystemExceptions": 0,
        "latestSystemExceptions": [],
        "numSinkExceptions": 0,
        "latestSinkExceptions": [],
        "numWrittenToSink": 0,
        "lastReceivedTime": 0,
        "workerId": "c-standalone-fw-localhost-8080"
      }
    }
  ]
}
but I don't know why my produced message is not reaching my api server
I have cross checked and the api is correct and available
d
wireshark to look at the network traffic would be your best bet
o
do you think by not including the authentication header in my config file might have caused this ?
Ok I think I found the error, I was using a wrong topic
👍 1
d
So everything is working now? 🤞
o
yes boss, just one last clarification
I am working on a project that requires every event or data to render accordingly so I opt in to use pulsar to queue them all
now I have succeeded in the sink implementation, my question now is this
is it possible to make a http request to pulsar sending an external data to it to process then deliver via the sink I just implemented
from all I have read you can only produce message or data to pulsar using wss:// or pulsar:// , so making request to the source or pulsar with http:// is not possible ?
d
There is an REST client available. Would that work?
o
yeah I know of this, the app I am talking is a RocketChat(RC) app, so it seems RC does not accept neither ws nor pulsar protocol just their inbuilt http even axios and fetch does not work
d
Maybe Pulsar Beam would be a better option?
🙌 1
o
thank you so very much boss, I will check it out
but this means there is no way I can directly make http request to pulsar from my webapp to sink my data ?
d
You can use whatever library you like to post to the HTTP endpoint exposed by BEAM
o
Thank you so much sir
Good day @David K please is it possible to produce message for a particular subscriber in a topic where there are so many other subscribers ? if yes please how can I achieve this
d
The pulsar protocol does’t have built-in support for this behavior. However, there are two ways to emulate this behavior. The first is to use Pulsar message properties to tag the message with its intended recipient’s ID. Consumers would then consume all the messages, but only process those that are intended for it. The downsides to this approach are that it is inefficient (all messages get sent over the network when only a fraction of them are processed), and that the data is exposed to all consumers, so sensitive data cannot be sent this way. You can solve the data security issue by using per-message encryption to control how can decrypt the message payload.
The second approach is to use some sort of server-side filtering mechanism that consumes messages from topics on the broker side, and then reroutes them to different topics behind the scenes based on message content and/or properties. This solves the inefficiency problem, but it requires the consumers to “know” ahead of time which topics are their’s to consume from. This is challenging in a scenario where the consumers are created/deleted dynamically such as the request-reply pattern.
o
the very first approach sounds great, what flag can I use to add this properties ?
d
I would recommend using the message properties, which is just a key/value map. Then pick a hard-coded key value to use for passing this information, e.g. “MSG_RECIPIENT” , etc.
o
oh thanks, the entire idea I am trying to implement is to control a certain customer's message production rate using
--dispatch-rate-period, -dt
and
--msg-dispatch-rate, -md
in a topic called
public/default/messenger
which has so many other subscribers, so I am just trying to figure out how I can do this with pulsar
any clue ?
d
The only back pressure mechanism that Pulsar has to control/throttle producers are backlog quotas, but those are at the topic level and not the producer level.
So if you created a separate topic for each customer, you could use backlog quotas to throttle them. Then you feed messages from those topics into a central topic for processing etc.
o
ok, one more question please
why does pulsar not allow me to use my self created tenant, namespace for sink processes example for the http-sink I am using to run my app, it only allows the public tenant and default namespace
d
It should allow you to use whatever tenant/namespace you like. W/o seeing the error, it sounds like a permissions issue. Did you grant yourself produce/consume and function permissions ?
o
no I didn’t, when is this permission set ?
d
It is set at the namespace or topic level. What is the error you are seeing when you try to submit the sink
o
please what is it that I am doing wrong here
Copy code
const axios = require('axios');
const FormData = require('form-data');
let data = new FormData();
data.append('classname', 'org.example.MySinkTest');
data.append('sourceSubscriptionName', 'Latest');
data.append('inputs', '["<persistent://public/default/ojoxdan>"]');
data.append('topicsPattern', 'public/default/');
data.append('inputSpecs', '{"schemaType": "type-x", "serdeClassName": "name-x", "isRegexPattern": true, "receiverQueueSize": 5}');
data.append('configs', '{"url":"<http://localhost:9000/consume/sms/smsenvio>"}');
data.append('secrets', 'hello');
data.append('parallelism', '1');
data.append('processingGuarantees', 'EFFECTIVELY_ONCE');
data.append('retainOrdering', 'true');
data.append('resources', '{"cpu":1,"ram":2,"disk":3}');
data.append('autoAck', 'true');
data.append('timeoutMs', '1000');
data.append('cleanupSubscription', 'true');
data.append('runtimeFlags', 'hello');
data.append('archive', '<https://dlcdn.apache.org/pulsar/pulsar-3.0.0/connectors/pulsar-io-http-3.0.0.nar>');

let config = {
  method: 'post',
  maxBodyLength: Infinity,
  url: '<http://localhost:8080/admin/v3/sinks/public/default/msgs>',
  headers: { 
    'Content-Type': 'multipart/form-data', 
    ...data.getHeaders()
  },
  data : data
};

axios.request(config)
.then((response) => {
  console.log(JSON.stringify(response.data));
})
.catch((error) => {
  console.log(error);
});
no matter what I add to this, it always return this response
Copy code
{
  "reason": "Sink config is not provided"
}
a
@ojoxdan I’m not super familiar with js but I think it happens because your sinkConfig is not really passed as a
multipart/form-data
form param in the request
o
@Alexander Preuß I used postman but got the same result. I had the content type set to application/json and used formdata for the data
a
@ojoxdan this is how the from parameters need to be passed:
Copy code
curl --location '<http://localhost:8080/admin/v3/sinks/public/default/my-jdbc>' \
--header 'Content-Type: multipart/form-data' \
--form 'sinkConfig="{"tenant":"public", "namespace":"default", ... }"
o
@Alexander Preuß let me try this out asap because I need it like right now
a
The backend expects the sinkConfig to be a single form parameter. I believe what happens in your code is it creates a form parameter for every property individually (the documentation/swagger definition is also wrong on the pulsar site regarding this)
o
ok I am trying now
hello @Alexander Preuß I just tried with postman again and it persit
a
I think you are still missing a
}
which could cause a parsing error
o
here again
here is the json object
Copy code
{
  "classname": "org.apache.pulsar.io.http.HttpSink",
  "archive": "<https://dlcdn.apache.org/pulsar/pulsar-3.0.0/connectors/pulsar-io-http-3.0.0.nar>",
  "namespace": "default",
  "tenant": "public",
  "inputs": [
    "<persistent://public/default/ojoxdan>"
  ],
  "configs": {
    "url": "<http://localhost:9000/consume/sms/smsenvio>"
  }
}
hello @Alexander Preuß I think I am resolving it now, I just got a new error
Copy code
reason: 'Sink package is not provided'
a
@ojoxdan I believe you should be able to fix this one by removing the
classname
config propertiy
o
ok let me try
@Alexander Preuß it's same
@David K and @Alexander Preuß any further help on this error
Copy code
{ reason: 'Sink package is not provided' }
?
d
Are there any errors in the broker/function worker logs around the time of that you post this request?
o
no sure
let me check
Copy code
Details = tenant: "public"
namespace: "default"
name: "customer-95fff9756d2c922b3c407b1401dc14c6"
className: "org.apache.pulsar.functions.api.utils.IdentityFunction"
autoAck: true
parallelism: 1
source {
  typeClassName: "org.apache.pulsar.client.api.schema.GenericObject"
  inputSpecs {
    key: "fe7c681397aad382017803677645e370"
    value {
    }
  }
  cleanupSubscription: true
}
sink {
  className: "org.apache.pulsar.io.http.HttpSink"
  configs: "{\"url\":\"<http://localhost:9000/consume/sms/smsenvio>\"}"
  typeClassName: "org.apache.pulsar.client.api.schema.GenericObject"
}
resources {
  cpu: 1.0
  ram: 1073741824
  disk: 10737418240
}
componentType: SINK
@David K here is what I found
Copy code
[public/default/customer-95fff9756d2c922b3c407b1401dc14c6-0] ERROR org.apache.pulsar.functions.instance.JavaInstanceRunnable - [public/default/customer-95fff9756d2c922b3c407b1401dc14c6:0] Uncaught exception in Java Instance
java.lang.NoClassDefFoundError: org/apache/pulsar/io/http/HttpSinkConfig
	at org.apache.pulsar.io.http.HttpSink.open(HttpSink.java:61) ~[?:?]
d
It looks like a packaging issue. If you examine the nar file contents, do you see the HttpSinkConfig class?
you might want to build the NAR file yourself from the source? Where did you download that NAR file from?
o
You know I have been trying to use the httpSink connector, I was able to perform the overall flow seamlessly via sh, but the approach is too slow, when I tried the api approach it was a speed of light compared to the cmd (sh) way, so now I am trying to rewrite the entire script
🤔 1
@David K any luck ?
d
Sorry, any luck with what? Rebuilding the NAR to include the missing class? Fixing the slowness of the httpSink in general?
o
I downloaded the file from the official repo
ok, let ask directly how can I create a httpSink using apache pulsar api approach
I need to provide the archive, and config url
lolz @David K same me on this hahahaha