[Apache Camel 2.8.0, JAXB, Jboss 6]
Przykład zaczytanie pliku XML do obiektu JAVA, z walidacją względem pliku XSD.
Wymagane dependency:
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-jaxb</artifactId>
<version>${camel-version}</version>
</dependency>
//W serwisie dodajemy metodę do uruchomienia pobierania danych z XML:
public void getBookXML() {
log.info("Import z XML BEGIN");
BookXML bookXML = null;
Exchange exchFile;
List<BookDTO> books = new ArrayList<BookDTO>();
ProducerTemplate producer = camelContext.createProducerTemplate();
ConsumerTemplate consumer = camelContext.createConsumerTemplate();
try {
String endpoint = MyUtils.getParameterOrNull("endpoint.file");
log.info("endpoint: "+endpoint);
exchFile = consumer.receive(endpoint, 2000L);
if (exchFile != null && exchFile.getIn() != null
&& exchFile.getIn().getBody() != null) {
log.info("Pobrano z: "
+ exchFile.getIn().getBody().getClass().toString());
log.info("ret: " + exchFile.getIn().getBody().toString());
Object validRet = producer.sendBody("direct:unmarshal",
ExchangePattern.InOut, exchFile.getIn().getBody());
if(validRet instanceof Boolean
&& (Boolean)validRet){
bookXML = (BookXML) producer.sendBody("direct:readXmlBook",
ExchangePattern.InOut, exchFile.getIn().getBody());
}else{
throw new MyException("B\u0142\u0105d walidacji XML");
}
books = prepareHashMapBooks(bookXML);
consumer.doneUoW(exchFile);
} else {
log.info("Brak pliku do importu");
}
log.info("Pobrano: "+books.size());
} catch (Exception e) {
e.printStackTrace();
}
log.info("Import z XML END");
}
//BUILDER TRAS APACHE CAMEL
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.model.dataformat.JaxbDataFormat;
public class MyRouteBuilderFile extends RouteBuilder{
@Override
public void configure() throws Exception {
JaxbDataFormat dataFormat = new JaxbDataFormat(false);
dataFormat.setContextPath("com.bloger.programmingmt.example");
//przeksztalcenie na z XML na DTO
from("direct:unmarshalFile")
.unmarshal(dataFormat)
.to("log: IN ==>");
//Parsowania XML na podstawie XSD
from("direct:unmarshal")
.convertBodyTo(org.w3c.dom.Document.class)
.bean(ParserXMLFileProcessor .class, "process")
.end();
}
}
//PROCESOR WALIDUJACY STRUKTURE PLIKU XML NA PODSTAWIE XSD
public class ParserXMLFileProcessor implements Processor {
private static final Logger log = Logger.getLogger(ParserXMLFileProcessor.class);
@Autowired
private CamelContext camelContext;
@Override
public void process(Exchange exchange) throws Exception {
Exchange exchFile = null;
try {
SchemaFactory factory = SchemaFactory
.newInstance(XMLConstants.W3C_XML_SCHEMA_NS_URI);
ConsumerTemplate consumer = camelContext.createConsumerTemplate();
//Wskazanie pliku XSD do walidacji
String schemaFileEndpoint = Utils.getParameterOrNull("endpoint.book.file.validator");
log.info("schemaFileEndpoint: "+schemaFileEndpoint);
try {
if(schemaFileEndpoint!=null
&& schemaFileEndpoint.length()>0){
exchFile = consumer.receive(schemaFileEndpoint, 2000L);
}
} catch (Exception e) {
e.printStackTrace();
throw new MyException("Wystąpił wyjątek podczas pobierania pliku XSD dla walidacji: "+e.getMessage());
}
if (exchFile != null && exchFile.getIn() != null
&& exchFile.getIn().getBody() != null) {
log.info("Pobrano z: "
+ exchFile.getIn().getBody().getClass().toString());
log.info("ret: " + exchFile.getIn().getBody().toString());
Schema schema = factory.newSchema((File)((GenericFile)exchFile.getIn().getBody()).getFile());
Validator validator = schema.newValidator();
validator.validate(new DOMSource((Document)exchange.getIn().getBody()));
System.out.println("...passed...");
consumer.doneUoW(exchFile);
} else {
log.info("Brak pliku do walidacji struktury");
}
exchange.getIn().setBody(Boolean.TRUE);
exchange.getOut().setBody(Boolean.TRUE);
} catch (Exception e) {
e.printStackTrace();
exchange.getIn().setBody(Boolean.FALSE);
exchange.getOut().setBody(Boolean.FALSE);
}
}
public void setCamelContext(
CamelContext camelContext) {
this.camelContext = camelContext;
}
}
Pokazywanie postów oznaczonych etykietą Apache Camel. Pokaż wszystkie posty
Pokazywanie postów oznaczonych etykietą Apache Camel. Pokaż wszystkie posty
poniedziałek, 22 lipca 2013
czwartek, 28 czerwca 2012
Apache Camel - klastrowanie Quartz.
[Apache Camel 2.8.0, Quartz 1.8.4, Jboss 6]
Konfiguracja Quartz w Spring dla środowiska klastrowego.
Wymagane zależności Maven:
<dependency>
<groupId>org.quartz-scheduler</groupId>
<artifactId>quartz-oracle</artifactId>
<version>1.8.4</version>
</dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-quartz</artifactId>
<version>2.8.0</version>
</dependency>
W pliku Context.xml umieszczamy:
<bean id="quartz" class="org.apache.camel.component.quartz.QuartzComponent">
<property name="scheduler" ref="scheduler"/>
<property name="autoStartScheduler" value="true"/>
</bean>
<bean id="scheduler" class="org.springframework.scheduling.quartz.SchedulerFactoryBean">
<property name="overwriteExistingJobs" value="true"/>
<property name="applicationContextSchedulerContextKey" value="applicationContext" />
<property name="autoStartup" value="false"/>
<property name="schedulerContextAsMap">
<map>
<entry key="CamelQuartzCamelContext" value-ref="camelContextCore"/>
</map>
</property>
<property name="quartzProperties">
<props>
<prop key="org.quartz.scheduler.instanceName">${org.quartz.scheduler.instanceName}</prop>
<prop key="org.quartz.scheduler.instanceId">${org.quartz.scheduler.instanceId}</prop>
<prop key="org.quartz.scheduler.skipUpdateCheck">${org.quartz.scheduler.skipUpdateCheck}</prop>
<prop key="org.quartz.threadPool.class">${org.quartz.threadPool.class}</prop>
<prop key="org.quartz.threadPool.threadCount">${org.quartz.threadPool.threadCount}</prop>
<prop key="org.quartz.threadPool.threadPriority">${org.quartz.threadPool.threadPriority}</prop>
<prop key="org.quartz.jobStore.misfireThreshold">${org.quartz.jobStore.misfireThreshold}</prop>
<prop key="org.quartz.jobStore.class">${org.quartz.jobStore.class}</prop>
<prop key="org.quartz.jobStore.driverDelegateClass">${org.quartz.jobStore.driverDelegateClass}</prop>
<prop key="org.quartz.jobStore.useProperties">${org.quartz.jobStore.useProperties}</prop>
<prop key="org.quartz.jobStore.dataSource">${org.quartz.jobStore.dataSource}</prop>
<prop key="org.quartz.jobStore.tablePrefix">${org.quartz.jobStore.tablePrefix}</prop>
<prop key="org.quartz.jobStore.isClustered">${org.quartz.jobStore.isClustered}</prop>
<prop key="org.quartz.jobStore.clusterCheckinInterval">${org.quartz.jobStore.clusterCheckinInterval}</prop>
<prop key="org.quartz.dataSource.jobScheduler.jndiURL">${org.quartz.dataSource.jobScheduler.jndiURL}
</prop>
<prop key="org.quartz.dataSource.jobScheduler.validationQuery">
${org.quartz.dataSource.jobScheduler.validationQuery}
</prop>
</props>
</property>
</bean>
Plik quartz.properties powinien wyglądać następująco:
#============================================================================
# Configure Main Scheduler Properties
#============================================================================
org.quartz.scheduler.instanceName = ClusteredScheduler
org.quartz.scheduler.instanceId = AUTO
#============================================================================
# Configure ThreadPool
#============================================================================
org.quartz.threadPool.class = org.quartz.simpl.SimpleThreadPool
org.quartz.threadPool.threadCount = 25
org.quartz.threadPool.threadPriority = 5
#============================================================================
# Configure JobStore
#============================================================================
org.quartz.jobStore.misfireThreshold = 60000
org.quartz.jobStore.class = org.quartz.impl.jdbcjobstore.JobStoreTX
org.quartz.jobStore.driverDelegateClass = org.quartz.impl.jdbcjobstore.oracle.OracleDelegate
org.quartz.jobStore.useProperties = false
org.quartz.jobStore.dataSource = jobScheduler
org.quartz.jobStore.tablePrefix = MY_QRTZ_
org.quartz.jobStore.isClustered = true
org.quartz.jobStore.clusterCheckinInterval = 20000
#============================================================================
# Configure Datasources
#============================================================================
org.quartz.dataSource.jobScheduler.jndiURL = java:MyDS
org.quartz.dataSource.jobScheduler.validationQuery=select 0 from dual
#Konfiguracja do polaczenia bez uzycia jndiName
#org.quartz.dataSource.jobScheduler.driver = oracle.jdbc.driver.OracleDriver
#org.quartz.dataSource.jobScheduler.URL = jdbc:oracle:thin:@localhost:1521:ssid
#org.quartz.dataSource.jobScheduler.user = tp2
#org.quartz.dataSource.jobScheduler.password = tp2
#org.quartz.dataSource.jobScheduler.maxConnections = 5
#org.quartz.dataSource.jobScheduler.validationQuery=select 0 from dual
Skrypt SQL do wykonania na wskazanym schemacie bazy danych:
DELETE FROM MY_QRTZ_JOB_LISTENERS;
DELETE FROM MY_QRTZ_TRIGGER_LISTENERS;
DELETE FROM MY_QRTZ_FIRED_TRIGGERS;
DELETE FROM MY_QRTZ_SIMPLE_TRIGGERS;
DELETE FROM MY_QRTZ_CRON_TRIGGERS;
DELETE FROM MY_QRTZ_BLOB_TRIGGERS;
DELETE FROM MY_QRTZ_TRIGGERS;
DELETE FROM MY_QRTZ_JOB_DETAILS;
DELETE FROM MY_QRTZ_CALENDARS;
DELETE FROM MY_QRTZ_PAUSED_TRIGGER_GRPS;
DELETE FROM MY_QRTZ_LOCKS;
DELETE FROM MY_QRTZ_SCHEDULER_STATE;
DROP TABLE MY_QRTZ_CALENDARS;
DROP TABLE MY_QRTZ_FIRED_TRIGGERS;
DROP TABLE MY_QRTZ_TRIGGER_LISTENERS;
DROP TABLE MY_QRTZ_BLOB_TRIGGERS;
DROP TABLE MY_QRTZ_CRON_TRIGGERS;
DROP TABLE MY_QRTZ_SIMPLE_TRIGGERS;
DROP TABLE MY_QRTZ_TRIGGERS;
DROP TABLE MY_QRTZ_JOB_LISTENERS;
DROP TABLE MY_QRTZ_JOB_DETAILS;
DROP TABLE MY_QRTZ_PAUSED_TRIGGER_GRPS;
DROP TABLE MY_QRTZ_LOCKS;
DROP TABLE MY_QRTZ_SCHEDULER_STATE;
CREATE TABLE MY_QRTZ_JOB_DETAILS
(
JOB_NAME VARCHAR2(80) NOT NULL,
JOB_GROUP VARCHAR2(80) NOT NULL,
DESCRIPTION VARCHAR2(120) NULL,
JOB_CLASS_NAME VARCHAR2(128) NOT NULL,
IS_DURABLE VARCHAR2(1) NOT NULL,
IS_VOLATILE VARCHAR2(1) NOT NULL,
IS_STATEFUL VARCHAR2(1) NOT NULL,
REQUESTS_RECOVERY VARCHAR2(1) NOT NULL,
JOB_DATA BLOB NULL,
PRIMARY KEY (JOB_NAME,JOB_GROUP)
);
CREATE TABLE MY_QRTZ_JOB_LISTENERS
(
JOB_NAME VARCHAR2(80) NOT NULL,
JOB_GROUP VARCHAR2(80) NOT NULL,
JOB_LISTENER VARCHAR2(80) NOT NULL,
PRIMARY KEY (JOB_NAME,JOB_GROUP,JOB_LISTENER),
FOREIGN KEY (JOB_NAME,JOB_GROUP)
REFERENCES MY_QRTZ_JOB_DETAILS(JOB_NAME,JOB_GROUP)
);
CREATE TABLE MY_QRTZ_TRIGGERS
(
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
JOB_NAME VARCHAR2(80) NOT NULL,
JOB_GROUP VARCHAR2(80) NOT NULL,
IS_VOLATILE VARCHAR2(1) NOT NULL,
DESCRIPTION VARCHAR2(120) NULL,
NEXT_FIRE_TIME NUMBER(13) NULL,
PREV_FIRE_TIME NUMBER(13) NULL,
TRIGGER_STATE VARCHAR2(16) NOT NULL,
TRIGGER_TYPE VARCHAR2(8) NOT NULL,
START_TIME NUMBER(13) NOT NULL,
END_TIME NUMBER(13) NULL,
CALENDAR_NAME VARCHAR2(80) NULL,
MISFIRE_INSTR NUMBER(2) NULL,
JOB_DATA BLOB NULL,
PRIORITY NUMBER(13) NULL,
PRIMARY KEY (TRIGGER_NAME,TRIGGER_GROUP),
FOREIGN KEY (JOB_NAME,JOB_GROUP)
REFERENCES QRTZ_JOB_DETAILS(JOB_NAME,JOB_GROUP)
);
CREATE TABLE MY_QRTZ_SIMPLE_TRIGGERS
(
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
REPEAT_COUNT NUMBER(7) NOT NULL,
REPEAT_INTERVAL NUMBER(12) NOT NULL,
TIMES_TRIGGERED NUMBER(7) NOT NULL,
PRIMARY KEY (TRIGGER_NAME,TRIGGER_GROUP),
FOREIGN KEY (TRIGGER_NAME,TRIGGER_GROUP)
REFERENCES MY_QRTZ_TRIGGERS(TRIGGER_NAME,TRIGGER_GROUP)
);
CREATE TABLE MY_QRTZ_CRON_TRIGGERS
(
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
CRON_EXPRESSION VARCHAR2(80) NOT NULL,
TIME_ZONE_ID VARCHAR2(80),
PRIMARY KEY (TRIGGER_NAME,TRIGGER_GROUP),
FOREIGN KEY (TRIGGER_NAME,TRIGGER_GROUP)
REFERENCES MY_QRTZ_TRIGGERS(TRIGGER_NAME,TRIGGER_GROUP)
);
CREATE TABLE MY_QRTZ_BLOB_TRIGGERS
(
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
BLOB_DATA BLOB NULL,
PRIMARY KEY (TRIGGER_NAME,TRIGGER_GROUP),
FOREIGN KEY (TRIGGER_NAME,TRIGGER_GROUP)
REFERENCES MY_QRTZ_TRIGGERS(TRIGGER_NAME,TRIGGER_GROUP)
);
CREATE TABLE MY_QRTZ_TRIGGER_LISTENERS
(
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
TRIGGER_LISTENER VARCHAR2(80) NOT NULL,
PRIMARY KEY (TRIGGER_NAME,TRIGGER_GROUP,TRIGGER_LISTENER),
FOREIGN KEY (TRIGGER_NAME,TRIGGER_GROUP)
REFERENCES MY_QRTZ_TRIGGERS(TRIGGER_NAME,TRIGGER_GROUP)
);
CREATE TABLE MY_QRTZ_CALENDARS
(
CALENDAR_NAME VARCHAR2(80) NOT NULL,
CALENDAR BLOB NOT NULL,
PRIMARY KEY (CALENDAR_NAME)
);
CREATE TABLE MY_QRTZ_PAUSED_TRIGGER_GRPS
(
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
PRIMARY KEY (TRIGGER_GROUP)
);
CREATE TABLE MY_QRTZ_FIRED_TRIGGERS
(
ENTRY_ID VARCHAR2(95) NOT NULL,
TRIGGER_NAME VARCHAR2(80) NOT NULL,
TRIGGER_GROUP VARCHAR2(80) NOT NULL,
IS_VOLATILE VARCHAR2(1) NOT NULL,
INSTANCE_NAME VARCHAR2(80) NOT NULL,
FIRED_TIME NUMBER(13) NOT NULL,
STATE VARCHAR2(16) NOT NULL,
JOB_NAME VARCHAR2(80) NULL,
JOB_GROUP VARCHAR2(80) NULL,
IS_STATEFUL VARCHAR2(1) NULL,
REQUESTS_RECOVERY VARCHAR2(1) NULL,
PRIORITY NUMBER(13) NULL,
PRIMARY KEY (ENTRY_ID)
);
CREATE TABLE MY_QRTZ_SCHEDULER_STATE
(
INSTANCE_NAME VARCHAR2(80) NOT NULL,
LAST_CHECKIN_TIME NUMBER(13) NOT NULL,
CHECKIN_INTERVAL NUMBER(13) NOT NULL,
RECOVERER VARCHAR2(80) NULL,
PRIMARY KEY (INSTANCE_NAME)
);
CREATE TABLE MY_QRTZ_LOCKS
(
LOCK_NAME VARCHAR2(40) NOT NULL,
PRIMARY KEY (LOCK_NAME)
);
INSERT INTO MY_QRTZ_LOCKS VALUES('TRIGGER_ACCESS');
INSERT INTO MY_QRTZ_LOCKS VALUES('JOB_ACCESS');
INSERT INTO MY_QRTZ_LOCKS VALUES('CALENDAR_ACCESS');
INSERT INTO MY_QRTZ_LOCKS VALUES('STATE_ACCESS');
INSERT INTO MY_QRTZ_LOCKS VALUES('MISFIRE_ACCESS');
Napotkane problemy:
Wygląda na to że w wersji Camel 2.8.0 jest problem z CronScheduledRoutePolicy dla klastrowania.
Przy restarcie serwera/re-deploy aplikacji pojawia się wyjątek: org.quartz.ObjectAlreadyExistsException: Unable to store Job with name: 'job-...' and group: 'jobGroup-...', because one already exists with this identification.
Niestety na tą chwilę nie znalazłem innego sposobu jak napisanie klasy zarządzającej wskazaną trasą. Start/Stop odbywa się na zasadzie dwóch oddzielnych Quartz które podaję jako parametr. Zapewne ktoś powie że to mało eleganckie ... - nie twierdzę że to szczyt profesjonalizmu. Jednak wybrałem tą opcję po przeanalizowaniu za i przeciw.
piątek, 15 czerwca 2012
Apache Camel - wysyłanie wiadomości e-mail smtp
[Apache Camel 2.8.0]
PRzykłąd procesora wysyłającego wiadomość e-mail z uzyciem ApacheCamel
@Component("MyMailingProcessor")
public class MyMailingProcessor implements Processor {
@Override
public void process(Exchange exchange) throws Exception {
String accountpassword = "accountpassword";
String accountuser = "accountuser";
//Po ':' nalezy podac port
String accountsmtp = "127.0.0.1";
ProducerTemplate producer = exchange.getContext().createProducerTemplate();
StringBuffer mailMessage = new StringBuffer();
mailMessage.append("Tresc wiadomosci e-mail.");
if(exchange.getIn()!=null){
mailMessage.append("\n\nObiekt: ");
mailMessage.append(exchange.getIn().getBody().toString());
}else{
mailMessage.append("\n\nObiekt: brak");
}
if(accountuser!=null && accountuser.length()>0
&& accountsmtp!=null && accountsmtp.length()>0){
Map<String, Object> map = new HashMap<String, Object>();
map.put("To", "userTo@domail.com");
map.put("From", "userfrom@domail.com");
map.put("Subject", "subject");
map.put(Exchange.CONTENT_TYPE, "text/plain; charset=UTF-8");
try {
String endpoint = "smtp://"+accountsmtp+"?username="+accountuser+"&password="+accountpassword;
producer.sendBodyAndHeaders(endpoint
, mailMessage.toString(), map);
} catch (Exception e) {
log.error("Nie udalo sie wyslac wiadomosci: "
+ e.getMessage());
}
}
}
}
PRzykłąd procesora wysyłającego wiadomość e-mail z uzyciem ApacheCamel
@Component("MyMailingProcessor")
public class MyMailingProcessor implements Processor {
@Override
public void process(Exchange exchange) throws Exception {
String accountpassword = "accountpassword";
String accountuser = "accountuser";
//Po ':' nalezy podac port
String accountsmtp = "127.0.0.1";
ProducerTemplate producer = exchange.getContext().createProducerTemplate();
StringBuffer mailMessage = new StringBuffer();
mailMessage.append("Tresc wiadomosci e-mail.");
if(exchange.getIn()!=null){
mailMessage.append("\n\nObiekt: ");
mailMessage.append(exchange.getIn().getBody().toString());
}else{
mailMessage.append("\n\nObiekt: brak");
}
if(accountuser!=null && accountuser.length()>0
&& accountsmtp!=null && accountsmtp.length()>0){
Map<String, Object> map = new HashMap<String, Object>();
map.put("To", "userTo@domail.com");
map.put("From", "userfrom@domail.com");
map.put("Subject", "subject");
map.put(Exchange.CONTENT_TYPE, "text/plain; charset=UTF-8");
try {
String endpoint = "smtp://"+accountsmtp+"?username="+accountuser+"&password="+accountpassword;
producer.sendBodyAndHeaders(endpoint
, mailMessage.toString(), map);
} catch (Exception e) {
log.error("Nie udalo sie wyslac wiadomosci: "
+ e.getMessage());
}
}
}
}
Subskrybuj:
Posty (Atom)