@ -1,15 +1,16 @@ |
|||||
<Project> |
<Project> |
||||
<PropertyGroup> |
<PropertyGroup> |
||||
<LangVersion>latest</LangVersion> |
<LangVersion>latest</LangVersion> |
||||
<Version>3.1.0</Version> |
<Version>3.1.0</Version> |
||||
<NoWarn>$(NoWarn);CS1591</NoWarn> |
<NoWarn>$(NoWarn);CS1591</NoWarn> |
||||
<PackageIconUrl>https://abp.io/assets/abp_nupkg.png</PackageIconUrl> |
<PackageIconUrl>https://abp.io/assets/abp_nupkg.png</PackageIconUrl> |
||||
<PackageProjectUrl>https://abp.io</PackageProjectUrl> |
<PackageProjectUrl>https://abp.io/</PackageProjectUrl> |
||||
<PackageLicenseUrl>https://github.com/abpframework/abp/blob/master/LICENSE</PackageLicenseUrl> |
<PackageLicenseFile>LICENSE.md</PackageLicenseFile> |
||||
<RepositoryType>git</RepositoryType> |
<RepositoryType>git</RepositoryType> |
||||
<RepositoryUrl>https://github.com/abpframework/abp/</RepositoryUrl> |
<RepositoryUrl>https://github.com/abpframework/abp/</RepositoryUrl> |
||||
</PropertyGroup> |
</PropertyGroup> |
||||
<ItemGroup> |
<ItemGroup> |
||||
<PackageReference Include="SourceLink.Create.CommandLine" Version="2.8.3" PrivateAssets="All" /> |
<PackageReference Include="SourceLink.Create.CommandLine" Version="2.8.3" PrivateAssets="All" /> |
||||
|
<None Include="LICENSE.md" Pack="true" PackagePath=""/> |
||||
</ItemGroup> |
</ItemGroup> |
||||
</Project> |
</Project> |
||||
@ -0,0 +1,405 @@ |
|||||
|
# Introducing the Angular Service Proxy Generation |
||||
|
|
||||
|
Angular Service Proxy System **generates TypeScript services and models** to consume your backend HTTP APIs developed using the ABP Framework. So, you **don't manually create** models for your server side DTOs and perform raw HTTP calls to the server. |
||||
|
|
||||
|
ABP Framework has introduced the **new** Angular Service Proxy Generation system with the version 3.1. While this feature was available since the [v2.3](https://blog.abp.io/abp/ABP-Framework-v2_3_0-Has-Been-Released), it was not well covering some scenarios, like inheritance and generic types and had some known problems. **With the v3.1, we've re-written** it using the [Angular Schematics](https://angular.io/guide/schematics) system. Now, it is much more stable and feature rich. |
||||
|
|
||||
|
This post introduces the service proxy generation system and highlights some important features. |
||||
|
|
||||
|
## Installation |
||||
|
|
||||
|
### ABP CLI |
||||
|
|
||||
|
You need to have the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) to use the system. So, install it if you haven't installed before: |
||||
|
|
||||
|
````bash |
||||
|
dotnet tool install -g Volo.Abp.Cli |
||||
|
```` |
||||
|
|
||||
|
If you already have installed it before, you can update to the latest version: |
||||
|
|
||||
|
````shell |
||||
|
dotnet tool update -g Volo.Abp.Cli |
||||
|
```` |
||||
|
|
||||
|
### Project Configuration |
||||
|
|
||||
|
> If you've created your project with version 3.1 or later, you can skip this part since it will be already installed in your solution. |
||||
|
|
||||
|
For a solution that was created before v3.1, follow the steps below to configure the angular application: |
||||
|
|
||||
|
* Add `@abp/ng.schematics` package to the `devDependencies` of the Angular project. Run the following command in the root folder of the angular application: |
||||
|
|
||||
|
````bash |
||||
|
npm install @abp/ng.schematics --save-dev |
||||
|
```` |
||||
|
|
||||
|
- Add `rootNamespace` entry into the `apis/default` section in the `/src/environments/environment.ts`, as shown below: |
||||
|
|
||||
|
```json |
||||
|
apis: { |
||||
|
default: { |
||||
|
... |
||||
|
rootNamespace: 'Acme.BookStore' |
||||
|
}, |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
`Acme.BookStore` should be replaced by the root namespace of your .NET project. This ensures to not create unnecessary nested folders while creating the service proxy code. This value is `AngularProxyDemo` for the example solution explained below. |
||||
|
|
||||
|
## Basic Usage |
||||
|
|
||||
|
### Project Creation |
||||
|
|
||||
|
> If you already have a solution, you can skip this section. |
||||
|
|
||||
|
You need to [create](https://abp.io/get-started) your solution with the Angular UI. You can use the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) to create a new solution: |
||||
|
|
||||
|
````bash |
||||
|
abp new AngularProxyDemo -u angular |
||||
|
```` |
||||
|
|
||||
|
#### Run the Application |
||||
|
|
||||
|
The backend application must be up and running to be able to use the service proxy code generation system. |
||||
|
|
||||
|
> See the [getting started](https://docs.abp.io/en/abp/latest/Getting-Started?UI=NG&DB=EF&Tiered=No) guide if you don't know details of creating and running the solution. |
||||
|
|
||||
|
### Backend |
||||
|
|
||||
|
Assume that we have an `IBookAppService` interface: |
||||
|
|
||||
|
````csharp |
||||
|
using System.Collections.Generic; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.Application.Services; |
||||
|
|
||||
|
namespace AngularProxyDemo.Books |
||||
|
{ |
||||
|
public interface IBookAppService : IApplicationService |
||||
|
{ |
||||
|
public Task<List<BookDto>> GetListAsync(); |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
That uses a `BookDto` defined as shown: |
||||
|
|
||||
|
```csharp |
||||
|
using System; |
||||
|
using Volo.Abp.Application.Dtos; |
||||
|
|
||||
|
namespace AngularProxyDemo.Books |
||||
|
{ |
||||
|
public class BookDto : EntityDto<Guid> |
||||
|
{ |
||||
|
public string Name { get; set; } |
||||
|
|
||||
|
public DateTime PublishDate { get; set; } |
||||
|
} |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
And implemented as the following: |
||||
|
|
||||
|
```csharp |
||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.Application.Services; |
||||
|
|
||||
|
namespace AngularProxyDemo.Books |
||||
|
{ |
||||
|
public class BookAppService : ApplicationService, IBookAppService |
||||
|
{ |
||||
|
public async Task<List<BookDto>> GetListAsync() |
||||
|
{ |
||||
|
//TODO: get books from a database... |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
It simply returns a list of books. You probably want to get the books from a database, but it doesn't matter for this article. |
||||
|
|
||||
|
### HTTP API |
||||
|
|
||||
|
Thanks to the [auto API controllers](https://docs.abp.io/en/abp/latest/API/Auto-API-Controllers) system of the ABP Framework, we don't have to develop API controllers manually. Just **run the backend (*HttpApi.Host*) application** that shows the [Swagger UI](https://swagger.io/tools/swagger-ui/) by default. You will see the **GET** API for the books: |
||||
|
|
||||
|
 |
||||
|
|
||||
|
### Service Proxy Generation |
||||
|
|
||||
|
Open a **command line** in the **root folder of the Angular application** and execute the following command: |
||||
|
|
||||
|
````bash |
||||
|
abp generate-proxy |
||||
|
```` |
||||
|
|
||||
|
It should produce an output like the following: |
||||
|
|
||||
|
````bash |
||||
|
CREATE src/app/shared/models/books/index.ts (142 bytes) |
||||
|
CREATE src/app/shared/services/books/book.service.ts (437 bytes) |
||||
|
... |
||||
|
```` |
||||
|
|
||||
|
> `generate-proxy` command can take some some optional parameters for advanced scenarios (like [modular development](https://docs.abp.io/en/abp/latest/Module-Development-Basics)). You can take a look at the [documentation](https://docs.abp.io/en/abp/latest/UI/Angular/Service-Proxies). |
||||
|
|
||||
|
It basically creates two files; |
||||
|
|
||||
|
src/app/shared/services/books/**book.service.ts**: This is the service that can be injected and used to get the list of books; |
||||
|
|
||||
|
````typescript |
||||
|
import { RestService } from '@abp/ng.core'; |
||||
|
import { Injectable } from '@angular/core'; |
||||
|
import type { BookDto } from '../../models/book'; |
||||
|
|
||||
|
@Injectable({ |
||||
|
providedIn: 'root', |
||||
|
}) |
||||
|
export class BookService { |
||||
|
apiName = 'Default'; |
||||
|
|
||||
|
getList = () => |
||||
|
this.restService.request<any, BookDto[]>({ |
||||
|
method: 'GET', |
||||
|
url: `/api/app/book`, |
||||
|
}, |
||||
|
{ apiName: this.apiName }); |
||||
|
|
||||
|
constructor(private restService: RestService) {} |
||||
|
} |
||||
|
|
||||
|
```` |
||||
|
|
||||
|
src/app/shared/models/books/**index.ts**: This file contains the modal classes corresponding to the DTOs defined in the server side; |
||||
|
|
||||
|
````typescript |
||||
|
import type { EntityDto } from '@abp/ng.core'; |
||||
|
|
||||
|
export interface BookDto extends EntityDto<string> { |
||||
|
name: string; |
||||
|
publishDate: string; |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
You can now inject the `BookService` into any Angular component and use the `getList()` method to get the list of books. |
||||
|
|
||||
|
### About the Generated Code |
||||
|
|
||||
|
The generated code is; |
||||
|
|
||||
|
* **Simple**: It is almost identical to the code if you've written it yourself. |
||||
|
* **Splitted**: Instead of a single, large file; |
||||
|
* It creates a separate `.ts` file for every backend **service**. **Model** (DTO) classes are also grouped per service. |
||||
|
* It understands the [modularity](https://docs.abp.io/en/abp/latest/Module-Development-Basics), so creates the services for your own **module** (or the module you've specified). |
||||
|
* **Object oriented**; |
||||
|
* Supports **inheritance** of server side DTOs and generates the code respecting to the inheritance structure. |
||||
|
* Supports **generic types**. |
||||
|
* Supports **re-using type definitions** across services and doesn't generate the same DTO multiple times. |
||||
|
* **Well-aligned to the backend**; |
||||
|
* Service **method signatures** match exactly with the services on the backend services. This is achieved by a special endpoint exposed by the ABP Framework that well defines the backend contracts. |
||||
|
* **Namespaces** are exactly matches to the backend services and DTOs. |
||||
|
* **Well-aligned with the ABP Framework**; |
||||
|
* Recognizes the **standard ABP Framework DTO types** (like `EntityDto`, `ListResultDto`... etc) and doesn't repeat these classes in the application code, but uses from the `@abp/ng.core` package. |
||||
|
* Uses the `RestService` defined by the `@abp/ng.core` package which simplifies the generated code, keeps it short and re-uses all the logics implemented by the `RestService` (including error handling, authorization token injection, using multiple server endpoints... etc). |
||||
|
|
||||
|
These are the main motivations behind the decision of creating a service proxy generation system, instead of using a pre-built tool like [NSWAG](https://github.com/RicoSuter/NSwag). |
||||
|
|
||||
|
## Other Examples |
||||
|
|
||||
|
Let me show you a few more examples. |
||||
|
|
||||
|
### Updating an Entity |
||||
|
|
||||
|
Assume that you added a new method to the server side application service, to update a book: |
||||
|
|
||||
|
```csharp |
||||
|
public Task<BookDto> UpdateAsync(Guid id, BookUpdateDto input); |
||||
|
``` |
||||
|
|
||||
|
`BookUpdateDto` is a simple class defined shown below: |
||||
|
|
||||
|
```csharp |
||||
|
using System; |
||||
|
|
||||
|
namespace AngularProxyDemo.Books |
||||
|
{ |
||||
|
public class BookUpdateDto |
||||
|
{ |
||||
|
public string Name { get; set; } |
||||
|
|
||||
|
public DateTime PublishDate { get; set; } |
||||
|
} |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
Let's re-run the `generate-proxy` command to see the result: |
||||
|
|
||||
|
```` |
||||
|
abp generate-proxy |
||||
|
```` |
||||
|
|
||||
|
The output of this command will be like the following: |
||||
|
|
||||
|
````bash |
||||
|
UPDATE src/app/shared/services/books/book.service.ts (660 bytes) |
||||
|
UPDATE src/app/shared/models/books/index.ts (217 bytes) |
||||
|
```` |
||||
|
|
||||
|
It tells us two files have been updated. Let's see the changes; |
||||
|
|
||||
|
**book.service.ts** |
||||
|
|
||||
|
````typescript |
||||
|
import { RestService } from '@abp/ng.core'; |
||||
|
import { Injectable } from '@angular/core'; |
||||
|
import type { BookDto, BookUpdateDto } from '../../models/books'; |
||||
|
|
||||
|
@Injectable({ |
||||
|
providedIn: 'root', |
||||
|
}) |
||||
|
export class BookService { |
||||
|
apiName = 'Default'; |
||||
|
|
||||
|
getList = () => |
||||
|
this.restService.request<any, BookDto[]>({ |
||||
|
method: 'GET', |
||||
|
url: `/api/app/book`, |
||||
|
}, |
||||
|
{ apiName: this.apiName }); |
||||
|
|
||||
|
update = (id: string, input: BookUpdateDto) => |
||||
|
this.restService.request<any, BookDto>({ |
||||
|
method: 'PUT', |
||||
|
url: `/api/app/book/${id}`, |
||||
|
body: input, |
||||
|
}, |
||||
|
{ apiName: this.apiName }); |
||||
|
|
||||
|
constructor(private restService: RestService) {} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
`update` function has been added to the `BookService` that gets an `id` and a `BookUpdateDto` as the parameters. |
||||
|
|
||||
|
**index.ts** |
||||
|
|
||||
|
````typescript |
||||
|
import type { EntityDto } from '@abp/ng.core'; |
||||
|
|
||||
|
export interface BookDto extends EntityDto<string> { |
||||
|
name: string; |
||||
|
publishDate: string; |
||||
|
} |
||||
|
|
||||
|
export interface BookUpdateDto { |
||||
|
name: string; |
||||
|
publishDate: string; |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
Added a new DTO class: `BookUpdateDto`. |
||||
|
|
||||
|
### Advanced Example |
||||
|
|
||||
|
In this example, I want to show a DTO structure using inheritance, generics, arrays and dictionaries. |
||||
|
|
||||
|
I've created an `IOrderAppService` as shown below: |
||||
|
|
||||
|
````csharp |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.Application.Services; |
||||
|
|
||||
|
namespace AngularProxyDemo.Orders |
||||
|
{ |
||||
|
public interface IOrderAppService : IApplicationService |
||||
|
{ |
||||
|
public Task CreateAsync(OrderCreateDto input); |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
`OrderCreateDto` and the related DTOs are as the followings; |
||||
|
|
||||
|
````csharp |
||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using Volo.Abp.Data; |
||||
|
|
||||
|
namespace AngularProxyDemo.Orders |
||||
|
{ |
||||
|
public class OrderCreateDto : IHasExtraProperties |
||||
|
{ |
||||
|
public Guid CustomerId { get; set; } |
||||
|
|
||||
|
public DateTime CreationTime { get; set; } |
||||
|
|
||||
|
//ARRAY of DTOs |
||||
|
public OrderDetailDto[] Details { get; set; } |
||||
|
|
||||
|
//DICTIONARY |
||||
|
public Dictionary<string, object> ExtraProperties { get; set; } |
||||
|
} |
||||
|
|
||||
|
public class OrderDetailDto : GenericDetailDto<int> //INHERIT from GENERIC |
||||
|
{ |
||||
|
public string Note { get; set; } |
||||
|
} |
||||
|
|
||||
|
//GENERIC class |
||||
|
public abstract class GenericDetailDto<TCount> |
||||
|
{ |
||||
|
public Guid ProductId { get; set; } |
||||
|
|
||||
|
public TCount Count { get; set; } |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
When I run the `abp generate-proxy` command again, I see two new files have been created. |
||||
|
|
||||
|
src/app/shared/services/orders/**order.service.ts** |
||||
|
|
||||
|
````typescript |
||||
|
import { RestService } from '@abp/ng.core'; |
||||
|
import { Injectable } from '@angular/core'; |
||||
|
import type { OrderCreateDto } from '../../models/orders'; |
||||
|
|
||||
|
@Injectable({ |
||||
|
providedIn: 'root', |
||||
|
}) |
||||
|
export class OrderService { |
||||
|
apiName = 'Default'; |
||||
|
|
||||
|
create = (input: OrderCreateDto) => |
||||
|
this.restService.request<any, void>({ |
||||
|
method: 'POST', |
||||
|
url: `/api/app/order`, |
||||
|
body: input, |
||||
|
}, |
||||
|
{ apiName: this.apiName }); |
||||
|
|
||||
|
constructor(private restService: RestService) {} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
src/app/shared/models/orders/**index.ts** |
||||
|
|
||||
|
````typescript |
||||
|
import type { GenericDetailDto } from '../../models/orders'; |
||||
|
|
||||
|
export interface OrderCreateDto { |
||||
|
customerId: string; |
||||
|
creationTime: string; |
||||
|
details: OrderDetailDto[]; |
||||
|
extraProperties: string | object; |
||||
|
} |
||||
|
|
||||
|
export interface OrderDetailDto extends GenericDetailDto<number> { |
||||
|
note: string; |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
NOTE: 3.1.0-rc2 was generating the code above, which is wrong. It will be fixed in the next RC versions. |
||||
|
After Width: | Height: | Size: 42 KiB |
|
After Width: | Height: | Size: 9.3 KiB |
|
Before Width: | Height: | Size: 123 KiB After Width: | Height: | Size: 123 KiB |
|
Before Width: | Height: | Size: 11 KiB After Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 8.4 KiB After Width: | Height: | Size: 8.4 KiB |
|
Before Width: | Height: | Size: 44 KiB After Width: | Height: | Size: 44 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 106 KiB After Width: | Height: | Size: 106 KiB |
|
Before Width: | Height: | Size: 7.9 KiB After Width: | Height: | Size: 7.9 KiB |
|
Before Width: | Height: | Size: 6.2 KiB After Width: | Height: | Size: 6.2 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 11 KiB After Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 8.5 KiB After Width: | Height: | Size: 8.5 KiB |
|
Before Width: | Height: | Size: 68 KiB After Width: | Height: | Size: 68 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 4.5 KiB After Width: | Height: | Size: 4.5 KiB |
@ -0,0 +1,49 @@ |
|||||
|
# How to add the user entity as a navigation property? |
||||
|
|
||||
|
In this post, I'll show you how to add the user as a navigation property in your new entity. |
||||
|
|
||||
|
To do this, open the ABP Suite. Create a new entity called `Note`. |
||||
|
|
||||
|
 |
||||
|
|
||||
|
Then add a string property called `Title`. |
||||
|
|
||||
|
 |
||||
|
|
||||
|
To be able to add a user navigation, we need to create a user DTO to map from entity. To do this, create a new folder called "Users" in `*.Application.Contracts` then add a new class called `AppUserDto` inherited from `IdentityUserDto`. |
||||
|
|
||||
|
 |
||||
|
|
||||
|
Create the mapping for `AppUserDto`. To do this, open `YourProjectApplicationAutoMapperProfile.cs` and add the below line: |
||||
|
|
||||
|
```csharp |
||||
|
CreateMap<AppUser, AppUserDto>().Ignore(x => x.ExtraProperties); |
||||
|
``` |
||||
|
|
||||
|
 |
||||
|
|
||||
|
Get back to ABP Suite, go to **Navigation Properties** tab. Click **Add Navigation Property** button. Browse `AppUser.cs` in `*.Domain\Users` folder. Then choose the `Name` item as display property. Browse `AppUserDto.cs` in `*.Contracts\Users` folder. Choose `Users` from Collection Names dropdown. |
||||
|
|
||||
|
 |
||||
|
|
||||
|
That's it! Click **Save and generate** button to create your page. You'll see the following page if there's everything goes well. |
||||
|
|
||||
|
 |
||||
|
|
||||
|
|
||||
|
|
||||
|
Note this example is implemented with ABP Commercial 3.1.0-rc.3. This is a RC version. If you want to install the CLI and Suite RC version follow the next steps: |
||||
|
|
||||
|
1- Uninstall the current version of the CLI and install the specific RC version: |
||||
|
|
||||
|
```bash |
||||
|
dotnet tool uninstall --global Volo.Abp.Cli && dotnet tool install --global Volo.Abp.Cli --version 3.1.0-rc.3 |
||||
|
``` |
||||
|
|
||||
|
2- Uninstall the current version of the Suite and install the specific RC version: |
||||
|
|
||||
|
```bash |
||||
|
dotnet tool uninstall --global Volo.Abp.Suite && dotnet tool install -g Volo.Abp.Suite --version 3.1.0-rc.3 --add-source https://nuget.abp.io/<YOUR-API-KEY>/v3/index.json |
||||
|
``` |
||||
|
|
||||
|
Don't forget to replace the `<YOUR-API-KEY>` with your own key! |
||||
|
After Width: | Height: | Size: 84 KiB |
|
After Width: | Height: | Size: 223 KiB |
|
After Width: | Height: | Size: 180 KiB |
|
After Width: | Height: | Size: 302 KiB |
|
After Width: | Height: | Size: 153 KiB |
|
After Width: | Height: | Size: 92 KiB |
@ -0,0 +1,167 @@ |
|||||
|
# Distributed Event Bus Kafka Integration |
||||
|
|
||||
|
> This document explains **how to configure the [Kafka](https://kafka.apache.org/)** as the distributed event bus provider. See the [distributed event bus document](Distributed-Event-Bus.md) to learn how to use the distributed event bus system |
||||
|
|
||||
|
## Installation |
||||
|
|
||||
|
Use the ABP CLI to add [Volo.Abp.EventBus.Kafka](https://www.nuget.org/packages/Volo.Abp.EventBus.Kafka) NuGet package to your project: |
||||
|
|
||||
|
* Install the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) if you haven't installed before. |
||||
|
* Open a command line (terminal) in the directory of the `.csproj` file you want to add the `Volo.Abp.EventBus.Kafka` package. |
||||
|
* Run `abp add-package Volo.Abp.EventBus.Kafka` command. |
||||
|
|
||||
|
If you want to do it manually, install the [Volo.Abp.EventBus.Kafka](https://www.nuget.org/packages/Volo.Abp.EventBus.Kafka) NuGet package to your project and add `[DependsOn(typeof(AbpEventBusKafkaModule))]` to the [ABP module](Module-Development-Basics.md) class inside your project. |
||||
|
|
||||
|
## Configuration |
||||
|
|
||||
|
You can configure using the standard [configuration system](Configuration.md), like using the `appsettings.json` file, or using the [options](Options.md) classes. |
||||
|
|
||||
|
### `appsettings.json` file configuration |
||||
|
|
||||
|
This is the simplest way to configure the Kafka settings. It is also very strong since you can use any other configuration source (like environment variables) that is [supported by the AspNet Core](https://docs.microsoft.com/en-us/aspnet/core/fundamentals/configuration/). |
||||
|
|
||||
|
**Example: The minimal configuration to connect to a local kafka server with default configurations** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "localhost:9092" |
||||
|
} |
||||
|
}, |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
* `MyGroupId` is the name of this application, which is used as the **GroupId** on the Kakfa. |
||||
|
* `MyTopicName` is the **topic name**. |
||||
|
|
||||
|
See [the Kafka document](https://docs.confluent.io/current/clients/confluent-kafka-dotnet/api/Confluent.Kafka.html) to understand these options better. |
||||
|
|
||||
|
#### Connections |
||||
|
|
||||
|
If you need to connect to another server than the localhost, you need to configure the connection properties. |
||||
|
|
||||
|
**Example: Specify the host name (as an IP address)** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092" |
||||
|
} |
||||
|
}, |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
Defining multiple connections is allowed. In this case, you can specify the connection that is used for the event bus. |
||||
|
|
||||
|
**Example: Declare two connections and use one of them for the event bus** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092" |
||||
|
}, |
||||
|
"SecondConnection": { |
||||
|
"BootstrapServers": "321.321.321.321:9092" |
||||
|
} |
||||
|
}, |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName", |
||||
|
"ConnectionName": "SecondConnection" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
This allows you to use multiple RabbitMQ server in your application, but select one of them for the event bus. |
||||
|
|
||||
|
You can use any of the [ClientConfig](https://docs.confluent.io/current/clients/confluent-kafka-dotnet/api/Confluent.Kafka.ClientConfig.html) properties as the connection properties. |
||||
|
|
||||
|
**Example: Specify the socket timeout** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092", |
||||
|
"SocketTimeoutMs": 60000 |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
### The Options Classes |
||||
|
|
||||
|
`AbpRabbitMqOptions` and `AbpRabbitMqEventBusOptions` classes can be used to configure the connection strings and event bus options for the RabbitMQ. |
||||
|
|
||||
|
You can configure this options inside the `ConfigureServices` of your [module](Module-Development-Basics.md). |
||||
|
|
||||
|
**Example: Configure the connection** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.Connections.Default.BootstrapServers = "123.123.123.123:9092"; |
||||
|
options.Connections.Default.SaslUsername = "user"; |
||||
|
options.Connections.Default.SaslPassword = "pwd"; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**Example: Configure the consumer config** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureConsumer = config => |
||||
|
{ |
||||
|
config.GroupId = "MyGroupId"; |
||||
|
config.EnableAutoCommit = false; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**Example: Configure the producer config** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureProducer = config => |
||||
|
{ |
||||
|
config.MessageTimeoutMs = 6000; |
||||
|
config.Acks = Acks.All; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**Example: Configure the topic specification** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureTopic = specification => |
||||
|
{ |
||||
|
specification.ReplicationFactor = 3; |
||||
|
specification.NumPartitions = 3; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
Using these options classes can be combined with the `appsettings.json` way. Configuring an option property in the code overrides the value in the configuration file. |
||||
@ -0,0 +1,109 @@ |
|||||
|
# Environment |
||||
|
|
||||
|
Every application needs some ** environment ** variables. In Angular world, this is usually managed by `environment.ts`, `environment.prod.ts` and so on. It is the same for ABP as well. |
||||
|
|
||||
|
Current `Environment` configuration holds sub config classes as follows: |
||||
|
|
||||
|
```typescript |
||||
|
export interface Environment { |
||||
|
apis: Apis; |
||||
|
application: Application; |
||||
|
oAuthConfig: AuthConfig; |
||||
|
production: boolean; |
||||
|
remoteEnv?: RemoteEnv; |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
## Apis |
||||
|
|
||||
|
```typescript |
||||
|
export interface Apis { |
||||
|
[key: string]: ApiConfig; |
||||
|
default: ApiConfig; |
||||
|
} |
||||
|
|
||||
|
export interface ApiConfig { |
||||
|
[key: string]: string; |
||||
|
rootNamespace?: string; |
||||
|
url: string; |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
Api config has to have a default config and it may have some additional ones for different modules. |
||||
|
I.e. you may want to connect to different Apis for different modules. |
||||
|
|
||||
|
Take a look at following example |
||||
|
|
||||
|
```json |
||||
|
{ |
||||
|
// ... |
||||
|
"apis": { |
||||
|
"default": { |
||||
|
"url": "https://localhost:8080", |
||||
|
}, |
||||
|
"AbpIdentity": { |
||||
|
"url": "https://localhost:9090", |
||||
|
} |
||||
|
}, |
||||
|
// ... |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
When an api from `AbpIdentity` is called, the request will be sent to `"https://localhost:9090"`. |
||||
|
Everything else will be sent to `"https://localhost:8080"` |
||||
|
|
||||
|
* `rootNamespace` **(new)** : Root namespace of the related API. e.g. Acme.BookStore |
||||
|
|
||||
|
## Application |
||||
|
|
||||
|
```typescript |
||||
|
export interface Application { |
||||
|
name: string; |
||||
|
baseUrl?: string; |
||||
|
logoUrl?: string; |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
* `name`: Name of the backend Application. It is also used by `logo.component` if `logoUrl` is not provided. |
||||
|
* `logoUrl`: Url of the application logo. It is used by `logo.component` |
||||
|
* `baseUrl`: [For detailed information](./Multi-Tenancy.md#domain-tenant-resolver) |
||||
|
|
||||
|
|
||||
|
## AuthConfig |
||||
|
|
||||
|
For authentication, we use angular-oauth2-oidc. Please check their [docs](https://github.com/manfredsteyer/angular-oauth2-oidc) out |
||||
|
|
||||
|
## RemoteEnvironment |
||||
|
|
||||
|
Some applications need to integrate an existing config into the `environment` used throughout the application. |
||||
|
Abp Framework supports this out of box. |
||||
|
|
||||
|
To integrate an existing config json into the `environment`, you need to set `remoteEnv` |
||||
|
|
||||
|
```typescript |
||||
|
export type customMergeFn = ( |
||||
|
localEnv: Partial<Config.Environment>, |
||||
|
remoteEnv: any, |
||||
|
) => Config.Environment; |
||||
|
|
||||
|
export interface RemoteEnv { |
||||
|
url: string; |
||||
|
mergeStrategy: 'deepmerge' | 'overwrite' | customMergeFn; |
||||
|
method?: string; |
||||
|
headers?: ABP.Dictionary<string>; |
||||
|
} |
||||
|
``` |
||||
|
|
||||
|
* `url` *: Required. The url to be used to retrieve environment config |
||||
|
* `mergeStrategy` *: Required. Defines how the local and the remote `environment` json will be merged |
||||
|
* `deepmerge`: Both local and remote `environment` json will be merged recursively. If both configs have same nested path, the remote `environment` will be prioritized. |
||||
|
* `overwrite`: Remote `environment` will be used and local environment will be ignored. |
||||
|
* `customMergeFn`: You can also provide your own merge function as shown in the example. It will take two parameters, `localEnv: Partial<Config.Environment>` and `remoteEnv` and it needs to return a `Config.Environment` object. |
||||
|
* `method`: HTTP method to be used when retrieving environment config. Default: `GET` |
||||
|
* `headers`: If extra headers are needed for the request, it can be set through this field. |
||||
|
|
||||
|
|
||||
|
## What's Next? |
||||
|
|
||||
|
* [Service Proxies](./Service-Proxies.md) |
||||
|
|
||||
@ -0,0 +1,163 @@ |
|||||
|
# 分布式事件总线Kafka集成 |
||||
|
|
||||
|
> 本文解释了**如何配置[Kafka](https://kafka.apache.org/)**做为分布式总线提供程序. 参阅[分布式事件总线文档](Distributed-Event-Bus.md)了解如何使用分布式事件总线系统. |
||||
|
|
||||
|
## 安装 |
||||
|
|
||||
|
使用ABP CLI添加[Volo.Abp.EventBus.Kafka[Volo.Abp.EventBus.Kafka](https://www.nuget.org/packages/Volo.Abp.EventBus.Kafka)NuGet包到你的项目: |
||||
|
|
||||
|
* 安装[ABP CLI](https://docs.abp.io/en/abp/latest/CLI),如果你还没有安装. |
||||
|
* 在你想要安装 `Volo.Abp.EventBus.Kafka` 包的 `.csproj` 文件目录打开命令行(终端). |
||||
|
* 运行 `abp add-package Volo.Abp.EventBus.Kafka` 命令. |
||||
|
|
||||
|
如果你想要手动安装,安装[Volo.Abp.EventBus.Kafka](https://www.nuget.org/packages/Volo.Abp.EventBus.Kafka) NuGet 包到你的项目然后添加 `[DependsOn(typeof(AbpEventBusKafkaModule))]` 到你的项目[模块](Module-Development-Basics.md)类. |
||||
|
|
||||
|
## 配置 |
||||
|
|
||||
|
可以使用配置使用标准的[配置系统](Configuration.md),如 `appsettings.json` 文件,或[选项](Options.md)类. |
||||
|
|
||||
|
### `appsettings.json` 文件配置 |
||||
|
|
||||
|
这是配置Kafka设置最简单的方法. 它也非常强大,因为你可以使用[由AspNet Core支持的](https://docs.microsoft.com/en-us/aspnet/core/fundamentals/configuration/)的任何其他配置源(如环境变量). |
||||
|
|
||||
|
**示例:最小化配置与默认配置连接到本地的Kafka服务器** |
||||
|
|
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
* `MyGroupId` 是应用程序的名称,用于Kafka的**GroupId**. |
||||
|
* `MyTopicName` 是**topic名称**. |
||||
|
|
||||
|
参阅[Kafka文档](https://docs.confluent.io/current/clients/confluent-kafka-dotnet/api/Confluent.Kafka.html)更好的了解这些选项. |
||||
|
|
||||
|
#### 连接 |
||||
|
|
||||
|
如果需要连接到本地主机以外的另一台服务器,需要配置连接属性. |
||||
|
|
||||
|
**示例: 指定主机名 (如IP地址)** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092" |
||||
|
} |
||||
|
}, |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
允许定义多个连接. 在这种情况下,你可以指定用于事件总线的连接. |
||||
|
|
||||
|
**示例: 声明两个连接并将其中一个用于事件总线** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092" |
||||
|
}, |
||||
|
"SecondConnection": { |
||||
|
"BootstrapServers": "321.321.321.321:9092" |
||||
|
} |
||||
|
}, |
||||
|
"EventBus": { |
||||
|
"GroupId": "MyGroupId", |
||||
|
"TopicName": "MyTopicName", |
||||
|
"ConnectionName": "SecondConnection" |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
这允许你可以在你的应用程序使用多个Kafka服务器,但将其中一个做为事件总线. |
||||
|
|
||||
|
你可以使用任何[ClientConfig](https://docs.confluent.io/current/clients/confluent-kafka-dotnet/api/Confluent.Kafka.ClientConfig.html)属性作为连接属性. |
||||
|
|
||||
|
**示例: 指定socket超时时间** |
||||
|
|
||||
|
````json |
||||
|
{ |
||||
|
"Kafka": { |
||||
|
"Connections": { |
||||
|
"Default": { |
||||
|
"BootstrapServers": "123.123.123.123:9092", |
||||
|
"SocketTimeoutMs": 60000 |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
```` |
||||
|
|
||||
|
### 选项类 |
||||
|
|
||||
|
`AbpKafkaOptions` 和 `AbpKafkaEventBusOptions` 类用于配置Kafka的连接字符串和事件总线选项. |
||||
|
|
||||
|
你可以在你的[模块](Module-Development-Basics.md)的 `ConfigureServices` 方法配置选项. |
||||
|
|
||||
|
**示例: 配置连接** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.Connections.Default.BootstrapServers = "123.123.123.123:9092"; |
||||
|
options.Connections.Default.SaslUsername = "user"; |
||||
|
options.Connections.Default.SaslPassword = "pwd"; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**示例: 配置 consumer config** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureConsumer = config => |
||||
|
{ |
||||
|
config.GroupId = "MyGroupId"; |
||||
|
config.EnableAutoCommit = false; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**示例: 配置 producer config** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureProducer = config => |
||||
|
{ |
||||
|
config.MessageTimeoutMs = 6000; |
||||
|
config.Acks = Acks.All; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
**示例: 配置 topic specification** |
||||
|
|
||||
|
````csharp |
||||
|
Configure<AbpKafkaOptions>(options => |
||||
|
{ |
||||
|
options.ConfigureTopic = specification => |
||||
|
{ |
||||
|
specification.ReplicationFactor = 3; |
||||
|
specification.NumPartitions = 3; |
||||
|
}; |
||||
|
}); |
||||
|
```` |
||||
|
|
||||
|
使用这些选项类可以与 `appsettings.json` 组合在一起. 在代码中配置选项属性会覆盖配置文件中的值. |
||||
@ -0,0 +1,13 @@ |
|||||
|
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
||||
|
using System.Collections.Generic; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.CropperJs |
||||
|
{ |
||||
|
public class CropperJsScriptContributor : BundleContributor |
||||
|
{ |
||||
|
public override void ConfigureBundle(BundleConfigurationContext context) |
||||
|
{ |
||||
|
context.Files.AddIfNotContains("/libs/cropperjs/js/cropper.min.js"); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,13 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.CropperJs |
||||
|
{ |
||||
|
public class CropperJsStyleContributor : BundleContributor |
||||
|
{ |
||||
|
public override void ConfigureBundle(BundleConfigurationContext context) |
||||
|
{ |
||||
|
context.Files.AddIfNotContains("/libs/cropperjs/css/cropper.min.css"); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,16 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
||||
|
using Volo.Abp.AspNetCore.Mvc.UI.Packages.JQuery; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.StarRatingSvg |
||||
|
{ |
||||
|
[DependsOn(typeof(JQueryScriptContributor))] |
||||
|
public class StarRatingSvgScriptContributor : BundleContributor |
||||
|
{ |
||||
|
public override void ConfigureBundle(BundleConfigurationContext context) |
||||
|
{ |
||||
|
context.Files.AddIfNotContains("/libs/star-rating-svg/js/jquery.star-rating-svg.min.js"); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,13 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.StarRatingSvg |
||||
|
{ |
||||
|
public class StarRatingSvgStyleContributor : BundleContributor |
||||
|
{ |
||||
|
public override void ConfigureBundle(BundleConfigurationContext context) |
||||
|
{ |
||||
|
context.Files.AddIfNotContains("/libs/star-rating-svg/css/star-rating-svg.css"); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -1,53 +1,11 @@ |
|||||
using System; |
|
||||
using System.IO; |
|
||||
using System.Text; |
|
||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp.Cli.Args; |
|
||||
using Volo.Abp.Cli.Utils; |
|
||||
using Volo.Abp.DependencyInjection; |
|
||||
|
|
||||
namespace Volo.Abp.Cli.Commands |
namespace Volo.Abp.Cli.Commands |
||||
{ |
{ |
||||
public class GenerateProxyCommand : IConsoleCommand, ITransientDependency |
public class GenerateProxyCommand : ProxyCommandBase |
||||
{ |
{ |
||||
public Task ExecuteAsync(CommandLineArgs commandLineArgs) |
public const string Name = "generate-proxy"; |
||||
{ |
|
||||
var angularPath = $"angular.json"; |
|
||||
if (!File.Exists(angularPath)) |
|
||||
{ |
|
||||
throw new CliUsageException( |
|
||||
"angular.json file not found. You must run this command in the angular folder." + |
|
||||
Environment.NewLine + Environment.NewLine + |
|
||||
GetUsageInfo() |
|
||||
); |
|
||||
} |
|
||||
|
|
||||
CmdHelper.RunCmd("npx ng g @abp/ng.schematics:proxy"); |
|
||||
|
|
||||
return Task.CompletedTask; |
|
||||
} |
|
||||
|
|
||||
public string GetUsageInfo() |
|
||||
{ |
|
||||
var sb = new StringBuilder(); |
|
||||
|
|
||||
sb.AppendLine(""); |
|
||||
sb.AppendLine("Usage:"); |
|
||||
sb.AppendLine(""); |
|
||||
sb.AppendLine(" abp generate-proxy"); |
|
||||
sb.AppendLine(""); |
|
||||
sb.AppendLine("Examples:"); |
|
||||
sb.AppendLine(""); |
|
||||
sb.AppendLine(" abp generate-proxy"); |
|
||||
sb.AppendLine(""); |
|
||||
sb.AppendLine("See the documentation for more info: https://docs.abp.io/en/abp/latest/CLI"); |
|
||||
|
|
||||
return sb.ToString(); |
protected override string CommandName => Name; |
||||
} |
|
||||
|
|
||||
public string GetShortDescription() |
protected override string SchematicsCommandName => "proxy-add"; |
||||
{ |
|
||||
return "Generates Angular service proxies and DTOs to consume HTTP APIs."; |
|
||||
} |
|
||||
} |
} |
||||
} |
} |
||||
|
|||||
@ -0,0 +1,144 @@ |
|||||
|
using System; |
||||
|
using System.IO; |
||||
|
using System.Text; |
||||
|
using System.Threading.Tasks; |
||||
|
using Newtonsoft.Json.Linq; |
||||
|
using Volo.Abp.Cli.Args; |
||||
|
using Volo.Abp.Cli.Utils; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.Cli.Commands |
||||
|
{ |
||||
|
public abstract class ProxyCommandBase : IConsoleCommand, ITransientDependency |
||||
|
{ |
||||
|
protected abstract string CommandName { get; } |
||||
|
|
||||
|
protected abstract string SchematicsCommandName { get; } |
||||
|
|
||||
|
public Task ExecuteAsync(CommandLineArgs commandLineArgs) |
||||
|
{ |
||||
|
CheckAngularJsonFile(); |
||||
|
CheckNgSchematics(); |
||||
|
|
||||
|
var prompt = commandLineArgs.Options.ContainsKey("p") || commandLineArgs.Options.ContainsKey("prompt"); |
||||
|
var defaultValue = prompt ? null : "__default"; |
||||
|
|
||||
|
var module = commandLineArgs.Options.GetOrNull(Options.Module.Short, Options.Module.Long) ?? defaultValue; |
||||
|
var source = commandLineArgs.Options.GetOrNull(Options.Source.Short, Options.Source.Long) ?? defaultValue; |
||||
|
var target = commandLineArgs.Options.GetOrNull(Options.Target.Short, Options.Target.Long) ?? defaultValue; |
||||
|
|
||||
|
var commandBuilder = new StringBuilder("npx ng g @abp/ng.schematics:" + SchematicsCommandName); |
||||
|
|
||||
|
if (module != null) |
||||
|
{ |
||||
|
commandBuilder.Append($" --module {module}"); |
||||
|
} |
||||
|
|
||||
|
if (source != null) |
||||
|
{ |
||||
|
commandBuilder.Append($" --source {source}"); |
||||
|
} |
||||
|
|
||||
|
if (target != null) |
||||
|
{ |
||||
|
commandBuilder.Append($" --target {target}"); |
||||
|
} |
||||
|
|
||||
|
CmdHelper.RunCmd(commandBuilder.ToString()); |
||||
|
|
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
|
||||
|
private void CheckNgSchematics() |
||||
|
{ |
||||
|
var packageJsonPath = $"package.json"; |
||||
|
|
||||
|
if (!File.Exists(packageJsonPath)) |
||||
|
{ |
||||
|
throw new CliUsageException( |
||||
|
"package.json file not found" + |
||||
|
Environment.NewLine + |
||||
|
GetUsageInfo() |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
var schematicsPackageNode = |
||||
|
(string) JObject.Parse(File.ReadAllText(packageJsonPath))["devDependencies"]?["@abp/ng.schematics"]; |
||||
|
|
||||
|
if (schematicsPackageNode == null) |
||||
|
{ |
||||
|
throw new CliUsageException( |
||||
|
"\"@abp/ng.schematics\" NPM package should be installed to the devDependencies before running this command!" + |
||||
|
Environment.NewLine + |
||||
|
GetUsageInfo() |
||||
|
); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void CheckAngularJsonFile() |
||||
|
{ |
||||
|
var angularPath = $"angular.json"; |
||||
|
if (!File.Exists(angularPath)) |
||||
|
{ |
||||
|
throw new CliUsageException( |
||||
|
"angular.json file not found. You must run this command in the angular folder." + |
||||
|
Environment.NewLine + Environment.NewLine + |
||||
|
GetUsageInfo() |
||||
|
); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public string GetUsageInfo() |
||||
|
{ |
||||
|
var sb = new StringBuilder(); |
||||
|
|
||||
|
sb.AppendLine(""); |
||||
|
sb.AppendLine("Usage:"); |
||||
|
sb.AppendLine(""); |
||||
|
sb.AppendLine($" abp {CommandName}"); |
||||
|
sb.AppendLine(""); |
||||
|
sb.AppendLine("Options:"); |
||||
|
sb.AppendLine(""); |
||||
|
sb.AppendLine("-m|--module <module-name> (default: 'app') The name of the backend module you wish to generate proxies for."); |
||||
|
sb.AppendLine("-s|--source <source-name> (default: 'defaultProject') Angular project name to resolve the root namespace & API definition URL from."); |
||||
|
sb.AppendLine("-t|--target <target-name> (default: 'defaultProject') Angular project name to place generated code in."); |
||||
|
sb.AppendLine("-p|--prompt Asks the options from the command line prompt (for the missing options)"); |
||||
|
sb.AppendLine(""); |
||||
|
sb.AppendLine("See the documentation for more info: https://docs.abp.io/en/abp/latest/CLI"); |
||||
|
|
||||
|
return sb.ToString(); |
||||
|
} |
||||
|
|
||||
|
public string GetShortDescription() |
||||
|
{ |
||||
|
return "Generates Angular service proxies and DTOs to consume HTTP APIs."; |
||||
|
} |
||||
|
|
||||
|
public static class Options |
||||
|
{ |
||||
|
public static class Module |
||||
|
{ |
||||
|
public const string Short = "m"; |
||||
|
public const string Long = "module"; |
||||
|
} |
||||
|
|
||||
|
public static class Source |
||||
|
{ |
||||
|
public const string Short = "s"; |
||||
|
public const string Long = "source"; |
||||
|
} |
||||
|
|
||||
|
public static class Target |
||||
|
{ |
||||
|
public const string Short = "t"; |
||||
|
public const string Long = "target"; |
||||
|
} |
||||
|
|
||||
|
public static class Prompt |
||||
|
{ |
||||
|
public const string Short = "p"; |
||||
|
public const string Long = "prompt"; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
namespace Volo.Abp.Cli.Commands |
||||
|
{ |
||||
|
public class RemoveProxyCommand : ProxyCommandBase |
||||
|
{ |
||||
|
public const string Name = "remove-proxy"; |
||||
|
|
||||
|
protected override string CommandName => Name; |
||||
|
|
||||
|
protected override string SchematicsCommandName => "proxy-remove"; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,16 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>netstandard2.0</TargetFramework> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.EventBus\Volo.Abp.EventBus.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Kafka\Volo.Abp.Kafka.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,27 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Kafka; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Kafka |
||||
|
{ |
||||
|
[DependsOn( |
||||
|
typeof(AbpEventBusModule), |
||||
|
typeof(AbpKafkaModule))] |
||||
|
public class AbpEventBusKafkaModule : AbpModule |
||||
|
{ |
||||
|
public override void ConfigureServices(ServiceConfigurationContext context) |
||||
|
{ |
||||
|
var configuration = context.Services.GetConfiguration(); |
||||
|
|
||||
|
Configure<AbpKafkaEventBusOptions>(configuration.GetSection("Kafka:EventBus")); |
||||
|
} |
||||
|
|
||||
|
public override void OnApplicationInitialization(ApplicationInitializationContext context) |
||||
|
{ |
||||
|
context |
||||
|
.ServiceProvider |
||||
|
.GetRequiredService<KafkaDistributedEventBus>() |
||||
|
.Initialize(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,12 @@ |
|||||
|
namespace Volo.Abp.EventBus.Kafka |
||||
|
{ |
||||
|
public class AbpKafkaEventBusOptions |
||||
|
{ |
||||
|
|
||||
|
public string ConnectionName { get; set; } |
||||
|
|
||||
|
public string TopicName { get; set; } |
||||
|
|
||||
|
public string GroupId { get; set; } |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,208 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
using System.Threading.Tasks; |
||||
|
using Confluent.Kafka; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.Kafka; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Threading; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Kafka |
||||
|
{ |
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
[ExposeServices(typeof(IDistributedEventBus), typeof(KafkaDistributedEventBus))] |
||||
|
public class KafkaDistributedEventBus : EventBusBase, IDistributedEventBus, ISingletonDependency |
||||
|
{ |
||||
|
protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } |
||||
|
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } |
||||
|
protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; } |
||||
|
protected IKafkaSerializer Serializer { get; } |
||||
|
protected IProducerPool ProducerPool { get; } |
||||
|
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; } |
||||
|
protected ConcurrentDictionary<string, Type> EventTypes { get; } |
||||
|
protected IKafkaMessageConsumer Consumer { get; private set; } |
||||
|
|
||||
|
public KafkaDistributedEventBus( |
||||
|
IServiceScopeFactory serviceScopeFactory, |
||||
|
ICurrentTenant currentTenant, |
||||
|
IOptions<AbpKafkaEventBusOptions> abpKafkaEventBusOptions, |
||||
|
IKafkaMessageConsumerFactory messageConsumerFactory, |
||||
|
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, |
||||
|
IKafkaSerializer serializer, |
||||
|
IProducerPool producerPool) |
||||
|
: base(serviceScopeFactory, currentTenant) |
||||
|
{ |
||||
|
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; |
||||
|
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; |
||||
|
MessageConsumerFactory = messageConsumerFactory; |
||||
|
Serializer = serializer; |
||||
|
ProducerPool = producerPool; |
||||
|
|
||||
|
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>(); |
||||
|
EventTypes = new ConcurrentDictionary<string, Type>(); |
||||
|
} |
||||
|
|
||||
|
public void Initialize() |
||||
|
{ |
||||
|
Consumer = MessageConsumerFactory.Create( |
||||
|
AbpKafkaEventBusOptions.TopicName, |
||||
|
AbpKafkaEventBusOptions.GroupId, |
||||
|
AbpKafkaEventBusOptions.ConnectionName); |
||||
|
|
||||
|
Consumer.OnMessageReceived(ProcessEventAsync); |
||||
|
|
||||
|
SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); |
||||
|
} |
||||
|
|
||||
|
private async Task ProcessEventAsync(Message<string, byte[]> message) |
||||
|
{ |
||||
|
var eventName = message.Key; |
||||
|
var eventType = EventTypes.GetOrDefault(eventName); |
||||
|
if (eventType == null) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var eventData = Serializer.Deserialize(message.Value, eventType); |
||||
|
|
||||
|
await TriggerHandlersAsync(eventType, eventData); |
||||
|
} |
||||
|
|
||||
|
public IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class |
||||
|
{ |
||||
|
return Subscribe(typeof(TEvent), handler); |
||||
|
} |
||||
|
|
||||
|
public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) |
||||
|
{ |
||||
|
var handlerFactories = GetOrCreateHandlerFactories(eventType); |
||||
|
|
||||
|
if (factory.IsInFactories(handlerFactories)) |
||||
|
{ |
||||
|
return NullDisposable.Instance; |
||||
|
} |
||||
|
|
||||
|
handlerFactories.Add(factory); |
||||
|
|
||||
|
return new EventHandlerFactoryUnregistrar(this, eventType, factory); |
||||
|
} |
||||
|
|
||||
|
/// <inheritdoc/>
|
||||
|
public override void Unsubscribe<TEvent>(Func<TEvent, Task> action) |
||||
|
{ |
||||
|
Check.NotNull(action, nameof(action)); |
||||
|
|
||||
|
GetOrCreateHandlerFactories(typeof(TEvent)) |
||||
|
.Locking(factories => |
||||
|
{ |
||||
|
factories.RemoveAll( |
||||
|
factory => |
||||
|
{ |
||||
|
var singleInstanceFactory = factory as SingleInstanceHandlerFactory; |
||||
|
if (singleInstanceFactory == null) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
var actionHandler = singleInstanceFactory.HandlerInstance as ActionEventHandler<TEvent>; |
||||
|
if (actionHandler == null) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
return actionHandler.Action == action; |
||||
|
}); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
/// <inheritdoc/>
|
||||
|
public override void Unsubscribe(Type eventType, IEventHandler handler) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType) |
||||
|
.Locking(factories => |
||||
|
{ |
||||
|
factories.RemoveAll( |
||||
|
factory => |
||||
|
factory is SingleInstanceHandlerFactory handlerFactory && |
||||
|
handlerFactory.HandlerInstance == handler |
||||
|
); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
/// <inheritdoc/>
|
||||
|
public override void Unsubscribe(Type eventType, IEventHandlerFactory factory) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory)); |
||||
|
} |
||||
|
|
||||
|
/// <inheritdoc/>
|
||||
|
public override void UnsubscribeAll(Type eventType) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); |
||||
|
} |
||||
|
|
||||
|
public override async Task PublishAsync(Type eventType, object eventData) |
||||
|
{ |
||||
|
var eventName = EventNameAttribute.GetNameOrDefault(eventType); |
||||
|
var body = Serializer.Serialize(eventData); |
||||
|
|
||||
|
var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); |
||||
|
|
||||
|
await producer.ProduceAsync( |
||||
|
AbpKafkaEventBusOptions.TopicName, |
||||
|
new Message<string, byte[]> |
||||
|
{ |
||||
|
Key = eventName, Value = body |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType) |
||||
|
{ |
||||
|
return HandlerFactories.GetOrAdd( |
||||
|
eventType, |
||||
|
type => |
||||
|
{ |
||||
|
var eventName = EventNameAttribute.GetNameOrDefault(type); |
||||
|
EventTypes[eventName] = type; |
||||
|
return new List<IEventHandlerFactory>(); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType) |
||||
|
{ |
||||
|
var handlerFactoryList = new List<EventTypeWithEventHandlerFactories>(); |
||||
|
|
||||
|
foreach (var handlerFactory in HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key)) |
||||
|
) |
||||
|
{ |
||||
|
handlerFactoryList.Add( |
||||
|
new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); |
||||
|
} |
||||
|
|
||||
|
return handlerFactoryList.ToArray(); |
||||
|
} |
||||
|
|
||||
|
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) |
||||
|
{ |
||||
|
//Should trigger same type
|
||||
|
if (handlerEventType == targetEventType) |
||||
|
{ |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
//Should trigger for inherited types
|
||||
|
if (handlerEventType.IsAssignableFrom(targetEventType)) |
||||
|
{ |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,17 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>netstandard2.0</TargetFramework> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<PackageReference Include="Confluent.Kafka" Version="1.5.0" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Json\Volo.Abp.Json.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Threading\Volo.Abp.Threading.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,31 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Json; |
||||
|
using Volo.Abp.Modularity; |
||||
|
using Volo.Abp.Threading; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
[DependsOn( |
||||
|
typeof(AbpJsonModule), |
||||
|
typeof(AbpThreadingModule) |
||||
|
)] |
||||
|
public class AbpKafkaModule : AbpModule |
||||
|
{ |
||||
|
public override void ConfigureServices(ServiceConfigurationContext context) |
||||
|
{ |
||||
|
var configuration = context.Services.GetConfiguration(); |
||||
|
Configure<AbpKafkaOptions>(configuration.GetSection("Kafka")); |
||||
|
} |
||||
|
|
||||
|
public override void OnApplicationShutdown(ApplicationShutdownContext context) |
||||
|
{ |
||||
|
context.ServiceProvider |
||||
|
.GetRequiredService<IConsumerPool>() |
||||
|
.Dispose(); |
||||
|
|
||||
|
context.ServiceProvider |
||||
|
.GetRequiredService<IProducerPool>() |
||||
|
.Dispose(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,22 @@ |
|||||
|
using System; |
||||
|
using Confluent.Kafka; |
||||
|
using Confluent.Kafka.Admin; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class AbpKafkaOptions |
||||
|
{ |
||||
|
public KafkaConnections Connections { get; } |
||||
|
|
||||
|
public Action<ProducerConfig> ConfigureProducer { get; set; } |
||||
|
|
||||
|
public Action<ConsumerConfig> ConfigureConsumer { get; set; } |
||||
|
|
||||
|
public Action<TopicSpecification> ConfigureTopic { get; set; } |
||||
|
|
||||
|
public AbpKafkaOptions() |
||||
|
{ |
||||
|
Connections = new KafkaConnections(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,109 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Diagnostics; |
||||
|
using System.Linq; |
||||
|
using Confluent.Kafka; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class ConsumerPool : IConsumerPool, ISingletonDependency |
||||
|
{ |
||||
|
protected AbpKafkaOptions Options { get; } |
||||
|
|
||||
|
protected ConcurrentDictionary<string, IConsumer<string, byte[]>> Consumers { get; } |
||||
|
|
||||
|
protected TimeSpan TotalDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(10); |
||||
|
|
||||
|
public ILogger<ConsumerPool> Logger { get; set; } |
||||
|
|
||||
|
private bool _isDisposed; |
||||
|
|
||||
|
public ConsumerPool(IOptions<AbpKafkaOptions> options) |
||||
|
{ |
||||
|
Options = options.Value; |
||||
|
|
||||
|
Consumers = new ConcurrentDictionary<string, IConsumer<string, byte[]>>(); |
||||
|
Logger = new NullLogger<ConsumerPool>(); |
||||
|
} |
||||
|
|
||||
|
public virtual IConsumer<string, byte[]> Get(string groupId, string connectionName = null) |
||||
|
{ |
||||
|
connectionName ??= KafkaConnections.DefaultConnectionName; |
||||
|
|
||||
|
return Consumers.GetOrAdd( |
||||
|
connectionName, connection => |
||||
|
{ |
||||
|
var config = new ConsumerConfig(Options.Connections.GetOrDefault(connection)) |
||||
|
{ |
||||
|
GroupId = groupId, |
||||
|
EnableAutoCommit = false |
||||
|
}; |
||||
|
|
||||
|
Options.ConfigureConsumer?.Invoke(config); |
||||
|
|
||||
|
return new ConsumerBuilder<string, byte[]>(config).Build(); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
public void Dispose() |
||||
|
{ |
||||
|
if (_isDisposed) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
_isDisposed = true; |
||||
|
|
||||
|
if (!Consumers.Any()) |
||||
|
{ |
||||
|
Logger.LogDebug($"Disposed consumer pool with no consumers in the pool."); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var poolDisposeStopwatch = Stopwatch.StartNew(); |
||||
|
|
||||
|
Logger.LogInformation($"Disposing consumer pool ({Consumers.Count} consumers)."); |
||||
|
|
||||
|
var remainingWaitDuration = TotalDisposeWaitDuration; |
||||
|
|
||||
|
foreach (var consumer in Consumers.Values) |
||||
|
{ |
||||
|
var poolItemDisposeStopwatch = Stopwatch.StartNew(); |
||||
|
|
||||
|
try |
||||
|
{ |
||||
|
consumer.Close(); |
||||
|
consumer.Dispose(); |
||||
|
} |
||||
|
catch |
||||
|
{ |
||||
|
} |
||||
|
|
||||
|
poolItemDisposeStopwatch.Stop(); |
||||
|
|
||||
|
remainingWaitDuration = remainingWaitDuration > poolItemDisposeStopwatch.Elapsed |
||||
|
? remainingWaitDuration.Subtract(poolItemDisposeStopwatch.Elapsed) |
||||
|
: TimeSpan.Zero; |
||||
|
} |
||||
|
|
||||
|
poolDisposeStopwatch.Stop(); |
||||
|
|
||||
|
Logger.LogInformation( |
||||
|
$"Disposed Kafka Consumer Pool ({Consumers.Count} consumers in {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms)."); |
||||
|
|
||||
|
if (poolDisposeStopwatch.Elapsed.TotalSeconds > 5.0) |
||||
|
{ |
||||
|
Logger.LogWarning( |
||||
|
$"Disposing Kafka Consumer Pool got time greather than expected: {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms."); |
||||
|
} |
||||
|
|
||||
|
Consumers.Clear(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,10 @@ |
|||||
|
using System; |
||||
|
using Confluent.Kafka; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public interface IConsumerPool : IDisposable |
||||
|
{ |
||||
|
IConsumer<string, byte[]> Get(string groupId, string connectionName = null); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using System; |
||||
|
using System.Threading.Tasks; |
||||
|
using Confluent.Kafka; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public interface IKafkaMessageConsumer |
||||
|
{ |
||||
|
void OnMessageReceived(Func<Message<string, byte[]>, Task> callback); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,19 @@ |
|||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public interface IKafkaMessageConsumerFactory |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// Creates a new <see cref="IKafkaMessageConsumer"/>.
|
||||
|
/// Avoid to create too many consumers since they are
|
||||
|
/// not disposed until end of the application.
|
||||
|
/// </summary>
|
||||
|
/// <param name="topicName"></param>
|
||||
|
/// <param name="groupId"></param>
|
||||
|
/// <param name="connectionName"></param>
|
||||
|
/// <returns></returns>
|
||||
|
IKafkaMessageConsumer Create( |
||||
|
string topicName, |
||||
|
string groupId, |
||||
|
string connectionName = null); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public interface IKafkaSerializer |
||||
|
{ |
||||
|
byte[] Serialize(object obj); |
||||
|
|
||||
|
object Deserialize(byte[] value, Type type); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,10 @@ |
|||||
|
using System; |
||||
|
using Confluent.Kafka; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public interface IProducerPool : IDisposable |
||||
|
{ |
||||
|
IProducer<string, byte[]> Get(string connectionName = null); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,35 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using Confluent.Kafka; |
||||
|
using JetBrains.Annotations; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
[Serializable] |
||||
|
public class KafkaConnections : Dictionary<string, ClientConfig> |
||||
|
{ |
||||
|
public const string DefaultConnectionName = "Default"; |
||||
|
|
||||
|
[NotNull] |
||||
|
public ClientConfig Default |
||||
|
{ |
||||
|
get => this[DefaultConnectionName]; |
||||
|
set => this[DefaultConnectionName] = Check.NotNull(value, nameof(value)); |
||||
|
} |
||||
|
|
||||
|
public KafkaConnections() |
||||
|
{ |
||||
|
Default = new ClientConfig(); |
||||
|
} |
||||
|
|
||||
|
public ClientConfig GetOrDefault(string connectionName) |
||||
|
{ |
||||
|
if (TryGetValue(connectionName, out var connectionFactory)) |
||||
|
{ |
||||
|
return connectionFactory; |
||||
|
} |
||||
|
|
||||
|
return Default; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,156 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Linq; |
||||
|
using System.Threading.Tasks; |
||||
|
using Confluent.Kafka; |
||||
|
using Confluent.Kafka.Admin; |
||||
|
using JetBrains.Annotations; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.ExceptionHandling; |
||||
|
using Volo.Abp.Threading; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class KafkaMessageConsumer : IKafkaMessageConsumer, ITransientDependency, IDisposable |
||||
|
{ |
||||
|
public ILogger<KafkaMessageConsumer> Logger { get; set; } |
||||
|
|
||||
|
protected IConsumerPool ConsumerPool { get; } |
||||
|
|
||||
|
protected IExceptionNotifier ExceptionNotifier { get; } |
||||
|
|
||||
|
protected AbpKafkaOptions Options { get; } |
||||
|
|
||||
|
protected ConcurrentBag<Func<Message<string, byte[]>, Task>> Callbacks { get; } |
||||
|
|
||||
|
protected IConsumer<string, byte[]> Consumer { get; private set; } |
||||
|
|
||||
|
protected string ConnectionName { get; private set; } |
||||
|
|
||||
|
protected string GroupId { get; private set; } |
||||
|
|
||||
|
protected string TopicName { get; private set; } |
||||
|
|
||||
|
public KafkaMessageConsumer( |
||||
|
IConsumerPool consumerPool, |
||||
|
IExceptionNotifier exceptionNotifier, |
||||
|
IOptions<AbpKafkaOptions> options) |
||||
|
{ |
||||
|
ConsumerPool = consumerPool; |
||||
|
ExceptionNotifier = exceptionNotifier; |
||||
|
Options = options.Value; |
||||
|
Logger = NullLogger<KafkaMessageConsumer>.Instance; |
||||
|
|
||||
|
Callbacks = new ConcurrentBag<Func<Message<string, byte[]>, Task>>(); |
||||
|
} |
||||
|
|
||||
|
public virtual void Initialize( |
||||
|
[NotNull] string topicName, |
||||
|
[NotNull] string groupId, |
||||
|
string connectionName = null) |
||||
|
{ |
||||
|
Check.NotNull(topicName, nameof(topicName)); |
||||
|
Check.NotNull(groupId, nameof(groupId)); |
||||
|
TopicName = topicName; |
||||
|
ConnectionName = connectionName ?? KafkaConnections.DefaultConnectionName; |
||||
|
GroupId = groupId; |
||||
|
|
||||
|
AsyncHelper.RunSync(CreateTopicAsync); |
||||
|
Consume(); |
||||
|
} |
||||
|
|
||||
|
public virtual void OnMessageReceived(Func<Message<string, byte[]>, Task> callback) |
||||
|
{ |
||||
|
Callbacks.Add(callback); |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task CreateTopicAsync() |
||||
|
{ |
||||
|
using (var adminClient = new AdminClientBuilder(Options.Connections.GetOrDefault(ConnectionName)).Build()) |
||||
|
{ |
||||
|
var topic = new TopicSpecification |
||||
|
{ |
||||
|
Name = TopicName, |
||||
|
NumPartitions = 1, |
||||
|
ReplicationFactor = 1 |
||||
|
}; |
||||
|
|
||||
|
Options.ConfigureTopic?.Invoke(topic); |
||||
|
|
||||
|
try |
||||
|
{ |
||||
|
await adminClient.CreateTopicsAsync(new[] {topic}); |
||||
|
} |
||||
|
catch (CreateTopicsException e) |
||||
|
{ |
||||
|
if (!e.Error.Reason.Contains($"Topic '{TopicName}' already exists")) |
||||
|
{ |
||||
|
throw; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected virtual void Consume() |
||||
|
{ |
||||
|
Consumer = ConsumerPool.Get(GroupId, ConnectionName); |
||||
|
|
||||
|
Task.Factory.StartNew(async () => |
||||
|
{ |
||||
|
Consumer.Subscribe(TopicName); |
||||
|
|
||||
|
while (true) |
||||
|
{ |
||||
|
try |
||||
|
{ |
||||
|
var consumeResult = Consumer.Consume(); |
||||
|
|
||||
|
if (consumeResult.IsPartitionEOF) |
||||
|
{ |
||||
|
continue; |
||||
|
} |
||||
|
|
||||
|
await HandleIncomingMessage(consumeResult); |
||||
|
} |
||||
|
catch (ConsumeException ex) |
||||
|
{ |
||||
|
Logger.LogException(ex, LogLevel.Warning); |
||||
|
AsyncHelper.RunSync(() => ExceptionNotifier.NotifyAsync(ex, logLevel: LogLevel.Warning)); |
||||
|
} |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task HandleIncomingMessage(ConsumeResult<string, byte[]> consumeResult) |
||||
|
{ |
||||
|
try |
||||
|
{ |
||||
|
foreach (var callback in Callbacks) |
||||
|
{ |
||||
|
await callback(consumeResult.Message); |
||||
|
} |
||||
|
|
||||
|
Consumer.Commit(consumeResult); |
||||
|
} |
||||
|
catch (Exception ex) |
||||
|
{ |
||||
|
Logger.LogException(ex); |
||||
|
await ExceptionNotifier.NotifyAsync(ex); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public virtual void Dispose() |
||||
|
{ |
||||
|
if (Consumer == null) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
Consumer.Close(); |
||||
|
Consumer.Dispose(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,32 @@ |
|||||
|
using System; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class KafkaMessageConsumerFactory : IKafkaMessageConsumerFactory, ISingletonDependency, IDisposable |
||||
|
{ |
||||
|
protected IServiceScope ServiceScope { get; } |
||||
|
|
||||
|
public KafkaMessageConsumerFactory(IServiceScopeFactory serviceScopeFactory) |
||||
|
{ |
||||
|
ServiceScope = serviceScopeFactory.CreateScope(); |
||||
|
} |
||||
|
|
||||
|
public IKafkaMessageConsumer Create( |
||||
|
string topicName, |
||||
|
string groupId, |
||||
|
string connectionName = null) |
||||
|
{ |
||||
|
var consumer = ServiceScope.ServiceProvider.GetRequiredService<KafkaMessageConsumer>(); |
||||
|
consumer.Initialize(topicName, groupId, connectionName); |
||||
|
return consumer; |
||||
|
} |
||||
|
|
||||
|
public void Dispose() |
||||
|
{ |
||||
|
ServiceScope?.Dispose(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,102 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Diagnostics; |
||||
|
using System.Linq; |
||||
|
using Confluent.Kafka; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class ProducerPool : IProducerPool, ISingletonDependency |
||||
|
{ |
||||
|
protected AbpKafkaOptions Options { get; } |
||||
|
|
||||
|
protected ConcurrentDictionary<string, IProducer<string, byte[]>> Producers { get; } |
||||
|
|
||||
|
protected TimeSpan TotalDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(10); |
||||
|
|
||||
|
public ILogger<ProducerPool> Logger { get; set; } |
||||
|
|
||||
|
private bool _isDisposed; |
||||
|
|
||||
|
public ProducerPool(IOptions<AbpKafkaOptions> options) |
||||
|
{ |
||||
|
Options = options.Value; |
||||
|
|
||||
|
Producers = new ConcurrentDictionary<string, IProducer<string, byte[]>>(); |
||||
|
Logger = new NullLogger<ProducerPool>(); |
||||
|
} |
||||
|
|
||||
|
public virtual IProducer<string, byte[]> Get(string connectionName = null) |
||||
|
{ |
||||
|
connectionName ??= KafkaConnections.DefaultConnectionName; |
||||
|
|
||||
|
return Producers.GetOrAdd( |
||||
|
connectionName, connection => |
||||
|
{ |
||||
|
var config = Options.Connections.GetOrDefault(connection); |
||||
|
|
||||
|
Options.ConfigureProducer?.Invoke(new ProducerConfig(config)); |
||||
|
|
||||
|
return new ProducerBuilder<string, byte[]>(config).Build(); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
public void Dispose() |
||||
|
{ |
||||
|
if (_isDisposed) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
_isDisposed = true; |
||||
|
|
||||
|
if (!Producers.Any()) |
||||
|
{ |
||||
|
Logger.LogDebug($"Disposed producer pool with no producers in the pool."); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var poolDisposeStopwatch = Stopwatch.StartNew(); |
||||
|
|
||||
|
Logger.LogInformation($"Disposing producer pool ({Producers.Count} producers)."); |
||||
|
|
||||
|
var remainingWaitDuration = TotalDisposeWaitDuration; |
||||
|
|
||||
|
foreach (var producer in Producers.Values) |
||||
|
{ |
||||
|
var poolItemDisposeStopwatch = Stopwatch.StartNew(); |
||||
|
|
||||
|
try |
||||
|
{ |
||||
|
producer.Dispose(); |
||||
|
} |
||||
|
catch |
||||
|
{ |
||||
|
} |
||||
|
|
||||
|
poolItemDisposeStopwatch.Stop(); |
||||
|
|
||||
|
remainingWaitDuration = remainingWaitDuration > poolItemDisposeStopwatch.Elapsed |
||||
|
? remainingWaitDuration.Subtract(poolItemDisposeStopwatch.Elapsed) |
||||
|
: TimeSpan.Zero; |
||||
|
} |
||||
|
|
||||
|
poolDisposeStopwatch.Stop(); |
||||
|
|
||||
|
Logger.LogInformation( |
||||
|
$"Disposed Kafka Producer Pool ({Producers.Count} producers in {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms)."); |
||||
|
|
||||
|
if (poolDisposeStopwatch.Elapsed.TotalSeconds > 5.0) |
||||
|
{ |
||||
|
Logger.LogWarning( |
||||
|
$"Disposing Kafka Producer Pool got time greather than expected: {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms."); |
||||
|
} |
||||
|
|
||||
|
Producers.Clear(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,27 @@ |
|||||
|
using System; |
||||
|
using System.Text; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Json; |
||||
|
|
||||
|
namespace Volo.Abp.Kafka |
||||
|
{ |
||||
|
public class Utf8JsonKafkaSerializer : IKafkaSerializer, ITransientDependency |
||||
|
{ |
||||
|
private readonly IJsonSerializer _jsonSerializer; |
||||
|
|
||||
|
public Utf8JsonKafkaSerializer(IJsonSerializer jsonSerializer) |
||||
|
{ |
||||
|
_jsonSerializer = jsonSerializer; |
||||
|
} |
||||
|
|
||||
|
public byte[] Serialize(object obj) |
||||
|
{ |
||||
|
return Encoding.UTF8.GetBytes(_jsonSerializer.Serialize(obj)); |
||||
|
} |
||||
|
|
||||
|
public object Deserialize(byte[] value, Type type) |
||||
|
{ |
||||
|
return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); |
||||
|
} |
||||
|
} |
||||
|
} |
||||