Search Options

Results per page
Sort
Preferred Languages
Advance

Results 1 - 3 of 3 for sendMessage (0.1 sec)

  1. internal/event/target/kafka.go

    	if target.producer == nil {
    		return store.ErrNotConnected
    	}
    	msg, err := target.toProducerMessage(eventData)
    	if err != nil {
    		return err
    	}
    	_, _, err = target.producer.SendMessage(msg)
    	return err
    }
    
    // sendMultiple sends multiple messages to the kafka.
    func (target *KafkaTarget) sendMultiple(events []event.Event) error {
    	if target.producer == nil {
    		return store.ErrNotConnected
    Registered: Sun Nov 03 19:28:11 UTC 2024
    - Last Modified: Fri Sep 06 23:06:30 UTC 2024
    - 13.6K bytes
    - Viewed (0)
  2. internal/logger/target/kafka/kafka.go

    	}
    	logJSON, err := json.Marshal(&entry)
    	if err != nil {
    		return err
    	}
    	msg := sarama.ProducerMessage{
    		Topic: h.kconfig.Topic,
    		Value: sarama.ByteEncoder(logJSON),
    	}
    	_, _, err = h.producer.SendMessage(&msg)
    	if err != nil {
    		atomic.StoreInt32(&h.status, statusOffline)
    	} else {
    		atomic.StoreInt32(&h.status, statusOnline)
    	}
    	return err
    }
    
    // Init initialize kafka target
    Registered: Sun Nov 03 19:28:11 UTC 2024
    - Last Modified: Fri Sep 06 23:06:30 UTC 2024
    - 10.2K bytes
    - Viewed (0)
  3. internal/s3select/message.go

    		strconv.FormatInt(bytesProcessed, 10) + `</BytesProcessed><BytesReturned>` +
    		strconv.FormatInt(bytesReturned, 10) + `</BytesReturned></Stats>`)
    	return genMessage(statsHeader, payload)
    }
    
    // endMessage - indicates that the request is complete, and no more messages will be sent.
    // You should not assume that the request is complete until the client receives an End message.
    //
    // Header specification:
    Registered: Sun Nov 03 19:28:11 UTC 2024
    - Last Modified: Tue Aug 30 15:26:43 UTC 2022
    - 15.2K bytes
    - Viewed (0)
Back to top