Slackbot
05/27/2022, 12:17 AMJohn Dimeo
05/27/2022, 12:51 AMElliot Huntington
05/27/2022, 1:11 AMprocess method each time a message is consumed from an input topic.
I guess what I'm looking for is some way to make the spring framework own/manage the lifecycle of that function. I would like to be able to provide a jar file to pulsar that, ideally, would be a spring boot application. And then pulsar could load a bean that is a org.apache.pulsar.functions.api.Function instance. I'm not sure exactly how this would work.
Here is an example of something that I'm currently trying to experiment with that I haven't been able to test quite yet. This is not the ideal example because ultimately, if this does work, pulsar would still be responsible for invoking the function's constructor rather than just loading the function from the spring application context:
Some service interface:
```public interface TransformerService {
String transform(String s);
}```Some service implementation:
```import org.springframework.stereotype.Service;
@Service
public class UpperCaseTransformerService implements TransformerService {
@Override
public String transform(String s) {
return s.toUpperCase();
}
}```
Some Function:
```import org.apache.pulsar.functions.api.Context;
import org.apache.pulsar.functions.api.Function;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.ConfigurableApplicationContext;
@SpringBootApplication
public class UppercaseTransformer implements Function<String, String> {
final ConfigurableApplicationContext ctx;
public UppercaseTransformer() {
ctx = SpringApplication.run(UppercaseTransformer.class);
}
public String process(String s, Context context) throws Exception {
final TransformerService springBean = ctx.getBean(TransformerService.class);
return springBean.transform(s);
}
}```Ideally I would want use start.spring.io to create a simple uber executable jar file that configures and launches the spring application. The function would be a spring bean, and pulsar would have the ability to load that spring bean from the application context.
Elliot Huntington
05/27/2022, 1:18 AM(This would allow the jar file to provide a function initializer that pulsar could use to get a handle on the function instead of creating it itself. Maybe there would also be a tear down method in this api that pulsar could invoke to gracefully shutoff the Function, or stop the app. One advantage of this approach would be to simplify pulsar's Function interface such that any java.util.Function could be used without loosing the ability to pass in a configuration object. I'm rather new to pulsar, but is it common for the Context object to have different values for each invocation of the function when a message is processed? Or would it make sense to just pass the context once during a function initialization, and then use the function for each message without passing the context?context) ->org.apache.pulsar.functions.api.Contextorg.apache.pulsar.functions.api.Function
John Dimeo
05/27/2022, 1:34 AMUppercaseTransformer and set all the other values available here: https://pulsar.apache.org/docs/pulsar-admin/#create-1
Spring could certainly tell Pulsar the input and output schemas/types and the function implementation class (and perhaps assume guess a name) but not much else besides that... unless you're saying all of that would come through Spring's configuration mechanisms?Elliot Huntington
05/27/2022, 3:11 AMJohn Dimeo
05/27/2022, 3:26 AMElliot Huntington
07/22/2022, 12:19 PMDevin G. Bost
07/22/2022, 12:25 PMElliot Huntington
07/22/2022, 12:31 PMElliot Huntington
07/22/2022, 12:33 PMElliot Huntington
07/22/2022, 12:43 PMExchangeRateProvider was an interface called PulsarFunctionProvider where PulsarFunction is an interface that specifies everything Pulsar internals need to properly configure/use a PulsarFunction.John Dimeo
07/22/2022, 12:45 PMpulsar-admin command to deploy. For example:
@Doc(name = "odf-finsummaryfn")
@Requires(@Doc(name = SupporterSummaryChange.TOPIC))
@Produces(@Doc(name = "odf-financial-summary", type = FinancialSummary.class))
@Parallelism(3)
public class FinancialSummaryFn implements Function<SupporterSummaryChange, FinancialSummary>, JSONUtils, IsDocumented {
....
@Doc is just a little annotation I wrote to track some documentation metadata on a class. I'm using it here to specify the function name, the topics it requires/is an input, the topic it produces/is the output, and the default/starting parallelism
I then do this kind of thing later:
String[] args = new String[] {
"functions", mode,
"--jar", "/tmp/" + jar + ".jar",
"--inputs", first(fnConfig.getInputs()),
"--classname", fnConfig.getClassName(),
"--name", fnConfig.getName(),
"--subs-position", pos.name(),
"--function-config-file", "/tmp/function-config.yaml",
"--parallelism", fnConfig.getParallelism().toString()
};
and use the annotations to create a FunctionConfig and use ProcessBuilder to call the function deploy command for meElliot Huntington
07/22/2022, 12:51 PMJohn Dimeo
07/22/2022, 12:53 PM@Parallelism
val instances = Optional.ofNullable(fn.getClass().getAnnotation(Parallelism.class)).map(Parallelism::value).orElse(1);Elliot Huntington
07/22/2022, 12:56 PMJohn Dimeo
07/22/2022, 1:04 PMFunction, manage all other config externally (today)
⢠Implement Function, annotate your class with some things, manage the rest through CLI or YAML (my set up)
⢠Implement a more complete Function interface that specifies the same things you can annotate, manage the rest through CLI or YAML (your suggestion)Elliot Huntington
09/21/2023, 12:33 PMFunction<I, O> and package that in a jar file.
Then there is the question, how does that function get deployed into Pulsar? We need to use the Pulsar Admin API. Based on the examples in the deploying functions documentation it looks like this:
$ bin/pulsar-admin functions create \
--jar <jarFile> \
--classname <classname> \
--inputs <inputs> \
--output <output>
The parameters --jar and --classname are the only parameters the admin api uses to instantiate the function. And Pulsar internals currently uses reflection to invoke a no arg constructor to instantiate an object of the type specified by the classname located within the jar.
The other parameters: --inputs, --outputs, etc... are what I think you referred to above as āenvironmentalā configuration. These have nothing to do with instantiating a function, but they are used for wiring up (configuring) a function.
Can you provide me with the names of the source code files within the Pulsar codebase that handles all this initialization and configuration? Iāll have a look at it to see if I can better understand.John Dimeo
12/15/2023, 7:23 PM--inputs is a mix of both: the topic name (if you really leverage tenants and namespaces and don't inject environmental information into your topic names) and schema are ontological and should be about instantiating the function. The function will break if the data coming in isn't the entity/data contract it expects. However, the tenant and namespace are environmental.
In my example above:
@Doc(name = "odf-finsummaryfn")
@Requires(@Doc(name = SupporterSummaryChange.TOPIC))
@Produces(@Doc(name = "odf-financial-summary", type = FinancialSummary.class))
It deliberately does not include things like CPU and RAM requirements or if this is the DEVINT or PROD cluster. But still things that are very helpful to define with the function implementation itself instead of always having to "remember" that function X consumes data contract Y during deployment.
So yes, we are conflating, because we must :-)